diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d66683b9a..f0325745a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -148,6 +148,10 @@ jobs: run: | go test -race -coverprofile=coverage.out -covermode=atomic ./... + # The Avro codec's registry handling only compiles with the kafka tag. + - name: Test the Avro codec + run: go test -tags kafka ./internal/infrastructure/codec/ + - name: Upload coverage uses: codecov/codecov-action@v4 with: diff --git a/Makefile b/Makefile index 15ef214ac..96da2830c 100644 --- a/Makefile +++ b/Makefile @@ -37,7 +37,7 @@ PROTOC_GEN_GO_GRPC_VERSION ?= v1.6.1 PROTO_DIR := internal/tasks/proto PROTO_GEN_FILES := $(PROTO_DIR)/tasks.pb.go -.PHONY: poollink-dev poollink-dev-down poollink-dev-reset setup-tools fmt lint check-migrations join-check split-cloud-check pages-check kafka-check proto check-proto \ +.PHONY: poollink-dev poollink-dev-down poollink-dev-reset setup-tools fmt schemas lint check-migrations join-check split-cloud-check pages-check kafka-check proto check-proto \ up upgrade claim doctor cli seed-demo seed seed-plan sandbox sandbox-seed sandbox-simulate reset logs status stop down test-seed \ restart restart-go restart-all infra infra-down app app-down app-logs \ backend forms forms-web consumer worker run dev tracking realtime web \ @@ -82,6 +82,11 @@ cli-check: fmt: gofmt -w ./cmd ./internal +# Record the bus schemas the next release publishes. CI refuses a change the +# registry would refuse, and a compatible one until it is recorded here. +schemas: + go test ./internal/app/eventschemas -run TestPublishedSchemasStayCompatible -update + lint: check-migrations join-check split-cloud-check pages-check check-dockerfiles ./scripts/check-forms-mirror.sh $(GO_BIN)/golangci-lint run --timeout=5m diff --git a/cmd/backend/main.go b/cmd/backend/main.go index 19e960b6f..3150151d4 100644 --- a/cmd/backend/main.go +++ b/cmd/backend/main.go @@ -56,6 +56,7 @@ import ( "github.com/warmbly/warmbly/internal/app/email" "github.com/warmbly/warmbly/internal/app/emailsend" emailverifyapp "github.com/warmbly/warmbly/internal/app/emailverify" + "github.com/warmbly/warmbly/internal/app/eventschemas" "github.com/warmbly/warmbly/internal/app/feature" "github.com/warmbly/warmbly/internal/app/fleet" "github.com/warmbly/warmbly/internal/app/fleetnode" @@ -149,6 +150,7 @@ import ( "github.com/warmbly/warmbly/internal/tasks" "github.com/warmbly/warmbly/internal/tasks/proto" "github.com/warmbly/warmbly/internal/tasksched" + "github.com/warmbly/warmbly/internal/version" "golang.org/x/oauth2" "golang.org/x/oauth2/google" ) @@ -596,6 +598,15 @@ func main() { errs.CaptureFatal(err) log.Fatal(err) } + // Register this release's bus schemas before any node publishes them. + go func() { + rctx, cancel := context.WithTimeout(ctx, time.Minute) + defer cancel() + if err := eventschemas.Register(rctx, codecImpl); err != nil { + errs.CaptureException(err) + log.Printf("event schemas: %v", err) + } + }() bus, err := eventbus.FromEnv(kafkaBootstrapServers, kafkaSaslConfig) if err != nil { @@ -1154,6 +1165,7 @@ func main() { GithubRepo: getenvDefault("RELEASES_GITHUB_REPO", "warmbly/warmbly"), WorkerImageRepo: getenvDefault("RELEASES_WORKER_IMAGE_REPO", "ghcr.io/warmbly/warmbly/worker"), GithubToken: os.Getenv("RELEASES_GITHUB_TOKEN"), + SchemaGate: eventschemas.Gate(codecImpl, version.Version), }, fleetSettingsRepo, ) diff --git a/docs/content/docs/development/architecture.mdx b/docs/content/docs/development/architecture.mdx index bf68694fb..5cb261191 100644 --- a/docs/content/docs/development/architecture.mdx +++ b/docs/content/docs/development/architecture.mdx @@ -169,6 +169,8 @@ An empty `desired_version` means "no opinion" and must never be read as "downgra The backend is deliberately excluded from this. It is what tells everyone else their version, and a self-update that goes wrong leaves nothing to recover with. +On an instance that encodes the bus with Avro, the fleet also waits for the backend. A booting backend registers every schema its release publishes, and the fleet only moves to a release once the backend runs it and the Schema Registry has accepted those schemas. A refusal holds the fleet on the release it has and shows in the Fleet section, instead of reaching a worker that would then fail every event it reports. + The release check runs once when the backend boots, and only with `RELEASES_ENABLED=true`: it is off by default, so a self-hosted fleet is never rolled onto a vendor image. An operator moves the fleet by hand from the admin panel's Fleet section or with `warmblyctl fleet version` and `warmblyctl fleet channel`, which set the channel or pin a tag directly. Configuration is env-driven (`RELEASES_ENABLED`, `RELEASES_GITHUB_REPO`, `RELEASES_WORKER_IMAGE_REPO`, `RELEASES_GITHUB_TOKEN`) so self-hosters can point at their own fork and registry. See `internal/app/releases/service.go`. diff --git a/docs/content/docs/development/configuration.mdx b/docs/content/docs/development/configuration.mdx index 935c90b85..fba4e2706 100644 --- a/docs/content/docs/development/configuration.mdx +++ b/docs/content/docs/development/configuration.mdx @@ -312,7 +312,7 @@ A worker's command topic is named after the node id issued when it joined, so th -`json` needs nothing and is the default. `avro` resolves every event against `SCHEMA_REGISTRY_URL` and ships only in the `-kafka` images. Producers and consumers have to agree on one: there is no in-band marker, so a consumer on the other codec cannot read what is already on the bus. Change it by draining the bus, not in place. `PUBSUB_ENABLED` must agree across backend, consumer and realtime. +`json` needs nothing and is the default. `avro` resolves every event against `SCHEMA_REGISTRY_URL` and ships only in the `-kafka` images. Producers and consumers have to agree on one: there is no in-band marker, so a consumer on the other codec cannot read what is already on the bus. Change it by draining the bus, not in place. The envelope subjects need `BACKWARD` compatibility, because a new event type adds a union branch; a publisher that meets `FORWARD` or `FULL` sets its subject to `BACKWARD` itself. `PUBSUB_ENABLED` must agree across backend, consumer and realtime. `KAFKA_TRACKING_TOPIC` is read by the Rust publisher and the Go subscriber. Override it in one place only and opens and clicks stop being consumed, with no error anywhere. diff --git a/docs/content/docs/development/events.mdx b/docs/content/docs/development/events.mdx index 5abe3cdc9..246ec5a03 100644 --- a/docs/content/docs/development/events.mdx +++ b/docs/content/docs/development/events.mdx @@ -20,13 +20,15 @@ Message encoding is orthogonal to transport, selected by `CODEC_PROVIDER`: Schemas are derived from the Go structs (`internal/models/event_schema.go`), and each field is given the Avro default matching its zero value. This is what makes adding a field safe: a reader on the new schema can still decode data written under the schema registered before it, filling the field it did not carry from the default, so the registry accepts the new version. A newly added field with no default is rejected under `BACKWARD` compatibility, and because the publisher registers before it serializes, the rejection stops every publish on that topic rather than degrading one field. -Two things follow. The registered document is the marshalled schema, not `Schema.String()`, which omits defaults and would throw the guarantee away. And the registry stays on `BACKWARD` rather than `FORWARD`: adding a field is safe under both, but adding a new event type adds a union branch, which `BACKWARD` accepts and `FORWARD` refuses. +Two things follow. The registered document is the marshalled schema, not `Schema.String()`, which omits defaults and would throw the guarantee away. And the registry stays on `BACKWARD` rather than `FORWARD`: adding a field is safe under both, but adding a new event type adds a union branch, which `BACKWARD` accepts and `FORWARD` refuses. When a registration is refused and the subject's effective level is `FORWARD` or `FULL` (or their transitive forms), the publisher sets that subject to `BACKWARD` and registers again, so the registry key needs permission to change a subject's compatibility. Any other level, and any other refusal, is left as it is. + +Nothing in this reaches production untested. `internal/app/eventschemas` lists every schema a release publishes and keeps the one the registry last accepted under `testdata/`. CI decodes the recorded schema with the new one, which is the registry's own `BACKWARD` check, and fails a change the registry would refuse. It also fails a compatible change that has not been recorded, so a schema change is always in the diff: run `make schemas` and commit the result. On release the backend registers every listed schema when it boots, and the fleet does not move to the release until that has succeeded. The Rust tracking service follows both switches. It speaks Kafka only when compiled with its `kafka` cargo feature, which is what the `tracking:*-kafka` image is, and it encodes with `CODEC_PROVIDER` on either transport. So a Kafka deployment on `CODEC_PROVIDER=json` needs no Schema Registry at all, and `avro` there is refused at boot without `SCHEMA_REGISTRY_URL` rather than failing at the first published event. -The consumer decodes `jobs.worker-events` and `tracking-events` with the same `CODEC_PROVIDER`, and worker envelopes cannot be Avro. That makes `json` the only value a working deployment uses, and it is why the tracking publisher has to read the setting rather than always writing Avro: a JSON consumer handed Avro drops every open and click with nothing but a deserialize warning. +The consumer decodes `jobs.worker-events` and `tracking-events` with the same `CODEC_PROVIDER`, which is why the tracking publisher has to read the setting rather than always writing Avro: a JSON consumer handed Avro drops every open and click with nothing but a deserialize warning. ## Topics diff --git a/internal/app/advisor/service.go b/internal/app/advisor/service.go index 347e5c0f0..0f8d25999 100644 --- a/internal/app/advisor/service.go +++ b/internal/app/advisor/service.go @@ -162,6 +162,10 @@ func WithCopyJudge(asker typesafe.Asker, cache repository.CopyJudgmentRepository // least important cards their rewrite, and the next run picks them up. const maxNarrationsPerRun = 12 +// narrationBudget bounds the completions of one run so the summary and run +// record still land inside the caller's deadline. +const narrationBudget = 60 * time.Second + func (s *service) Evaluate(ctx context.Context, orgID uuid.UUID, trigger string) (*models.AdvisorSummary, error) { settings, err := s.repo.GetSettings(ctx, orgID) if err != nil { @@ -219,7 +223,9 @@ func (s *service) Evaluate(ctx context.Context, orgID uuid.UUID, trigger string) // rather than reporting problems that no longer exist. fixed := s.autopilot(ctx, orgID, settings, stored) - narrated := s.narrateBatch(ctx, orgID, stored) + nctx, cancelNarrate := context.WithTimeout(ctx, narrationBudget) + narrated := s.narrateBatch(nctx, orgID, stored) + cancelNarrate() summary, err := s.repo.Summary(ctx, orgID) if err != nil { @@ -303,7 +309,7 @@ func (s *service) narrateBatch(ctx context.Context, orgID uuid.UUID, findings [] // Detect, so the cap spends the budget on what matters. count := 0 for _, f := range findings { - if count >= maxNarrationsPerRun { + if count >= maxNarrationsPerRun || ctx.Err() != nil { break } if f.Narrated { diff --git a/internal/app/consumer/event_new_email.go b/internal/app/consumer/event_new_email.go index df46505b7..d5b30bb5d 100644 --- a/internal/app/consumer/event_new_email.go +++ b/internal/app/consumer/event_new_email.go @@ -45,6 +45,7 @@ func (s *JobsService) ingestNewEmail(ctx context.Context, e *models.JobEventNewE log.Warn().Msg("NEW_EMAIL event without a message body, dropping") return nil } + e.Message.ValidText() warmupToken := warmupTokenFromMessage(e.Message) if warmupToken != "" { handled, err := s.handleWarmupEmail(ctx, e, warmupToken) diff --git a/internal/app/eventschemas/eventschemas.go b/internal/app/eventschemas/eventschemas.go new file mode 100644 index 000000000..6077f51b1 --- /dev/null +++ b/internal/app/eventschemas/eventschemas.go @@ -0,0 +1,67 @@ +// Package eventschemas registers a release's bus schemas from the control +// plane, which rolls out first, and holds the fleet off a release the schema +// registry refused, so a refusal never reaches a worker mid-send. +package eventschemas + +import ( + "context" + "fmt" + "strings" + + "github.com/hamba/avro/v2" + "github.com/warmbly/warmbly/internal/events" + "github.com/warmbly/warmbly/internal/infrastructure/codec" + "github.com/warmbly/warmbly/internal/infrastructure/kafka" + "github.com/warmbly/warmbly/internal/models" +) + +// WorkerCommands names every per-node command topic (w.). +const WorkerCommands = "w.*" + +// Published is every schema this build publishes, by topic. +func Published() map[string]avro.Schema { + return map[string]avro.Schema{ + kafka.TopicWorkerEvents: models.JobEvent{}.Schema(), + WorkerCommands: models.WorkerEvent{}.Schema(), + events.TopicEmailEvents: events.EmailSentEvent{}.Schema(), + events.TopicWarmupEvents: events.WarmupEmailSentEvent{}.Schema(), + } +} + +// Register registers every published schema when c resolves against a +// registry. A codec without one has nothing to refuse. +func Register(ctx context.Context, c codec.Codec) error { + r, ok := c.(codec.SchemaRegistrar) + if !ok { + return nil + } + return r.RegisterSchemas(ctx, Published()) +} + +// Gate answers whether the fleet may move to tag. On a registry, the control +// plane has to be running that release and its schemas have to register; +// running is this binary's stamped version, and a dev build is not held. +func Gate(c codec.Codec, running string) func(ctx context.Context, tag string) error { + return func(ctx context.Context, tag string) error { + if _, ok := c.(codec.SchemaRegistrar); !ok { + return nil + } + if release(running) && release(tag) && base(running) != base(tag) { + return fmt.Errorf("the control plane runs %s, so the fleet waits for it to run %s and register its schemas", running, tag) + } + if err := Register(ctx, c); err != nil { + return fmt.Errorf("schema registry refused this release's event schemas: %w", err) + } + return nil + } +} + +// release is a tagged version rather than a dev or unstamped build. +func release(v string) bool { + return len(v) > 1 && v[0] == 'v' && v[1] >= '0' && v[1] <= '9' +} + +// base drops the image variant, so v1.2.3-kafka and v1.2.3 are one release. +func base(v string) string { + return strings.TrimSuffix(v, "-kafka") +} diff --git a/internal/app/eventschemas/eventschemas_test.go b/internal/app/eventschemas/eventschemas_test.go new file mode 100644 index 000000000..e126b52e7 --- /dev/null +++ b/internal/app/eventschemas/eventschemas_test.go @@ -0,0 +1,121 @@ +package eventschemas + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "flag" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/hamba/avro/v2" + "github.com/warmbly/warmbly/internal/infrastructure/codec" + "github.com/warmbly/warmbly/internal/infrastructure/kafka" + "github.com/warmbly/warmbly/internal/models" +) + +var update = flag.Bool("update", false, "record the current schemas as the snapshot (make schemas)") + +// snapshotPath is the recorded schema for topic, as the registry last accepted it. +func snapshotPath(topic string) string { + name := topic + if topic == WorkerCommands { + name = "worker-commands" + } + return filepath.Join("testdata", name+".avsc") +} + +func document(t *testing.T, s avro.Schema) []byte { + t.Helper() + doc, err := models.SchemaDocument(s) + if err != nil { + t.Fatal(err) + } + var out bytes.Buffer + if err := json.Indent(&out, doc, "", " "); err != nil { + t.Fatal(err) + } + return append(out.Bytes(), '\n') +} + +// TestPublishedSchemasStayCompatible is the registry's check, run before merge: +// a reader on the new schema has to decode what the recorded one wrote. +func TestPublishedSchemasStayCompatible(t *testing.T) { + for topic, schema := range Published() { + t.Run(topic, func(t *testing.T) { + path := snapshotPath(topic) + current := document(t, schema) + recorded, err := os.ReadFile(path) + if err == nil { + old, perr := avro.Parse(string(recorded)) + if perr != nil { + t.Fatalf("%s does not parse: %v", path, perr) + } + if cerr := avro.NewSchemaCompatibility().Compatible(schema, old); cerr != nil { + t.Fatalf("the %s schema is not BACKWARD compatible with %s, so the registry would refuse it and every publish on the topic would stop: %v", topic, path, cerr) + } + } else if !os.IsNotExist(err) || !*update { + t.Fatalf("no recorded schema for %s; run make schemas: %v", topic, err) + } + if *update { + if err := os.WriteFile(path, current, 0o644); err != nil { + t.Fatal(err) + } + return + } + if !bytes.Equal(current, recorded) { + t.Fatalf("the %s schema changed compatibly; run make schemas to record it in %s", topic, path) + } + }) + } +} + +type registry struct { + codec.Codec + err error + registered map[string]avro.Schema +} + +func (r *registry) RegisterSchemas(_ context.Context, s map[string]avro.Schema) error { + r.registered = s + return r.err +} + +func TestGate(t *testing.T) { + ctx := context.Background() + for _, tc := range []struct { + name, running, tag string + err error + held bool + }{ + {"the control plane runs the release", "v1.2.3", "v1.2.3", nil, false}, + {"the kafka image is the same release", "v1.2.3-kafka", "v1.2.3", nil, false}, + {"the control plane is behind", "v1.2.3", "v1.2.4", nil, true}, + {"a dev build still registers", "dev-abc", "v1.2.4", nil, false}, + {"the registry refuses", "v1.2.3", "v1.2.3", errors.New("409"), true}, + } { + t.Run(tc.name, func(t *testing.T) { + r := ®istry{err: tc.err} + err := Gate(r, tc.running)(ctx, tc.tag) + if (err != nil) != tc.held { + t.Fatalf("held = %v, want %v", err, tc.held) + } + if !tc.held && len(r.registered) != len(Published()) { + t.Fatalf("registered %d schemas, want %d", len(r.registered), len(Published())) + } + }) + } + + if err := Gate(codec.NewJSON(), "v1.2.3")(ctx, "v9.9.9"); err != nil { + t.Fatalf("a codec with no registry held the fleet: %v", err) + } +} + +func TestWorkerCommandsNamesThePerNodeTopics(t *testing.T) { + if !strings.HasPrefix(kafka.GetWorkerTopic("node"), strings.TrimSuffix(WorkerCommands, "*")) { + t.Fatalf("WorkerCommands %q does not match w.", WorkerCommands) + } +} diff --git a/internal/app/eventschemas/testdata/email-events.avsc b/internal/app/eventschemas/testdata/email-events.avsc new file mode 100644 index 000000000..981e0a395 --- /dev/null +++ b/internal/app/eventschemas/testdata/email-events.avsc @@ -0,0 +1,59 @@ +{ + "fields": [ + { + "default": "", + "name": "event_type", + "type": "string" + }, + { + "default": "", + "name": "task_id", + "type": "string" + }, + { + "default": "", + "name": "account_id", + "type": "string" + }, + { + "default": "", + "name": "campaign_id", + "type": "string" + }, + { + "default": "", + "name": "contact_id", + "type": "string" + }, + { + "default": "", + "name": "sequence_id", + "type": "string" + }, + { + "default": "", + "name": "message_id", + "type": "string" + }, + { + "default": "", + "name": "recipient", + "type": "string" + }, + { + "default": "", + "name": "subject", + "type": "string" + }, + { + "default": 0, + "name": "sent_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + } + ], + "name": "warmbly.events.EmailSentEvent", + "type": "record" +} diff --git a/internal/app/eventschemas/testdata/jobs.worker-events.avsc b/internal/app/eventschemas/testdata/jobs.worker-events.avsc new file mode 100644 index 000000000..1426e51a4 --- /dev/null +++ b/internal/app/eventschemas/testdata/jobs.worker-events.avsc @@ -0,0 +1,1056 @@ +{ + "fields": [ + { + "name": "type", + "type": "string" + }, + { + "name": "body", + "type": [ + { + "fields": [ + { + "default": "", + "name": "task_id", + "type": "string" + }, + { + "default": "", + "name": "email_account_id", + "type": "string" + }, + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "error_code", + "type": "string" + }, + { + "default": "", + "name": "error_type", + "type": "string" + }, + { + "default": "", + "name": "resolve_method", + "type": "string" + }, + { + "default": "", + "name": "message", + "type": "string" + }, + { + "default": false, + "name": "user_visible", + "type": "boolean" + }, + { + "default": "", + "name": "user_title", + "type": "string" + }, + { + "default": "", + "name": "user_message", + "type": "string" + }, + { + "default": "", + "name": "action_required", + "type": "string" + }, + { + "default": 0, + "name": "timestamp", + "type": "long" + } + ], + "name": "warmbly.events.EmailErrorEvent", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "id", + "type": "string" + }, + { + "default": 0, + "name": "uid", + "type": "long" + }, + { + "default": "\u0000\u0000\u0000\u0000\u0000\u0000\u0000\u0000", + "name": "mod_seq", + "type": { + "name": "warmbly.events.uint64", + "size": 8, + "type": "fixed" + } + }, + { + "default": 0, + "name": "mailbox", + "type": "long" + }, + { + "default": "", + "name": "folder_path", + "type": "string" + }, + { + "default": "", + "name": "folder", + "type": "string" + }, + { + "default": [], + "name": "flags", + "type": { + "items": "string", + "type": "array" + } + } + ], + "name": "warmbly.events.JobEventEmailUpdate", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "id", + "type": "string" + }, + { + "default": [], + "name": "flags", + "type": { + "items": "string", + "type": "array" + } + } + ], + "name": "warmbly.events.JobEventFlags", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "folder", + "type": "string" + }, + { + "default": "", + "name": "delta_link", + "type": "string" + } + ], + "name": "warmbly.events.JobEventGraphDeltaUpdate", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "\u0000\u0000\u0000\u0000\u0000\u0000\u0000\u0000", + "name": "history_id", + "type": "warmbly.events.uint64" + } + ], + "name": "warmbly.events.JobEventHistoryIDUpdate", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "original_message_id", + "type": "string" + }, + { + "default": "", + "name": "failed_recipient", + "type": "string" + }, + { + "default": "", + "name": "reason", + "type": "string" + } + ], + "name": "warmbly.events.JobEventInboundBounce", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "original_message_id", + "type": "string" + }, + { + "default": "", + "name": "complained_recipient", + "type": "string" + }, + { + "default": "", + "name": "provider", + "type": "string" + } + ], + "name": "warmbly.events.JobEventInboundComplaint", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "mailbox", + "type": "string" + }, + { + "default": 0, + "name": "uid_validity", + "type": "long" + }, + { + "default": false, + "name": "skipped", + "type": "boolean" + } + ], + "name": "warmbly.events.JobEventMailboxDelete", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "from", + "type": "string" + }, + { + "default": "", + "name": "to", + "type": "string" + } + ], + "name": "warmbly.events.JobEventMailboxRename", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": null, + "name": "data", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "name", + "type": "string" + }, + { + "default": [], + "name": "attributes", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": 0, + "name": "uid_validity", + "type": "long" + }, + { + "default": "\u0000\u0000\u0000\u0000\u0000\u0000\u0000\u0000", + "name": "highestmodseq", + "type": "warmbly.events.uint64" + }, + { + "default": 0, + "name": "uid_next", + "type": "long" + }, + { + "default": "", + "name": "delim", + "type": "string" + }, + { + "default": 0, + "name": "messages", + "type": "long" + }, + { + "default": 0, + "name": "updated_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + } + ], + "name": "warmbly.events.Mailbox", + "type": "record" + } + ] + } + ], + "name": "warmbly.events.JobEventMailboxUpdate", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": null, + "name": "message", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": 0, + "name": "mailbox", + "type": "long" + }, + { + "default": "", + "name": "folder_path", + "type": "string" + }, + { + "default": "", + "name": "folder", + "type": "string" + }, + { + "default": "", + "name": "provider_folder", + "type": "string" + }, + { + "default": "", + "name": "thread_id", + "type": "string" + }, + { + "default": "", + "name": "message_id", + "type": "string" + }, + { + "default": "", + "name": "gmail_id", + "type": "string" + }, + { + "default": "", + "name": "parent_id", + "type": "string" + }, + { + "default": 0, + "name": "uid", + "type": "long" + }, + { + "default": "\u0000\u0000\u0000\u0000\u0000\u0000\u0000\u0000", + "name": "mod_seq", + "type": "warmbly.events.uint64" + }, + { + "default": [], + "name": "flags", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "bcc", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "cc", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "from_addr", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "in_reply_to", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "reply_to", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "to_addr", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": "", + "name": "subject", + "type": "string" + }, + { + "default": 0, + "name": "size", + "type": "long" + }, + { + "default": 0, + "name": "internal_date", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + }, + { + "default": 0, + "name": "sent_date", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + }, + { + "default": "", + "name": "snippet", + "type": "string" + }, + { + "default": false, + "name": "seen", + "type": "boolean" + }, + { + "default": "", + "name": "body_text", + "type": "string" + }, + { + "default": 0, + "name": "updated_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + }, + { + "default": 0, + "name": "created_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + } + ], + "name": "warmbly.events.EmailMessageStoreData", + "type": "record" + } + ] + }, + { + "default": "", + "name": "report_original_message_id", + "type": "string" + } + ], + "name": "warmbly.events.JobEventNewEmail", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "id", + "type": "string" + }, + { + "default": "", + "name": "skipped_folder", + "type": "string" + } + ], + "name": "warmbly.events.JobEventRemoveEmail", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": { + "backfill_completed_at": null, + "backfill_cursor": { + "folders": {}, + "page_token": "" + }, + "backfill_since": null, + "backfill_started_at": null, + "backfill_status": "", + "backfill_synced": 0, + "deferred": 0, + "folders_skipped_cap": 0, + "folders_skipped_conflict": 0, + "last_synced_at": null, + "throttle_reason": "", + "throttled_until": null + }, + "name": "state", + "type": { + "fields": [ + { + "default": "", + "name": "backfill_status", + "type": "string" + }, + { + "default": { + "folders": {}, + "page_token": "" + }, + "name": "backfill_cursor", + "type": { + "fields": [ + { + "default": "", + "name": "page_token", + "type": "string" + }, + { + "default": {}, + "name": "folders", + "type": { + "type": "map", + "values": { + "fields": [ + { + "default": "", + "name": "next", + "type": "string" + }, + { + "default": 0, + "name": "uid", + "type": "long" + }, + { + "default": false, + "name": "done", + "type": "boolean" + } + ], + "name": "warmbly.events.SyncFolderCursor", + "type": "record" + } + } + } + ], + "name": "warmbly.events.SyncCursor", + "type": "record" + } + }, + { + "default": 0, + "name": "backfill_synced", + "type": "int" + }, + { + "default": null, + "name": "backfill_since", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": null, + "name": "backfill_started_at", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": null, + "name": "backfill_completed_at", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": null, + "name": "throttled_until", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": "", + "name": "throttle_reason", + "type": "string" + }, + { + "default": 0, + "name": "deferred", + "type": "int" + }, + { + "default": 0, + "name": "folders_skipped_cap", + "type": "int" + }, + { + "default": 0, + "name": "folders_skipped_conflict", + "type": "int" + }, + { + "default": null, + "name": "last_synced_at", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + } + ], + "name": "warmbly.events.SyncState", + "type": "record" + } + } + ], + "name": "warmbly.events.JobEventSyncState", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "access_token", + "type": "string" + }, + { + "default": "", + "name": "refresh_token", + "type": "string" + }, + { + "default": 0, + "name": "expires_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + } + ], + "name": "warmbly.events.JobEventTokenUpdate", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "rfc_message_id", + "type": "string" + }, + { + "default": "", + "name": "outcome", + "type": "string" + }, + { + "default": false, + "name": "recheck", + "type": "boolean" + } + ], + "name": "warmbly.events.JobEventWarmupRemovalChecked", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "task_id", + "type": "string" + }, + { + "default": false, + "name": "success", + "type": "boolean" + }, + { + "default": "", + "name": "message_id", + "type": "string" + }, + { + "default": "", + "name": "provider_msg_id", + "type": "string" + }, + { + "default": "", + "name": "thread_id", + "type": "string" + }, + { + "default": 0, + "name": "sent_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + }, + { + "default": null, + "name": "error", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "code", + "type": "string" + }, + { + "default": "", + "name": "type", + "type": "string" + }, + { + "default": "", + "name": "message", + "type": "string" + }, + { + "default": "", + "name": "resolve_method", + "type": "string" + }, + { + "default": false, + "name": "user_visible", + "type": "boolean" + }, + { + "default": "", + "name": "user_title", + "type": "string" + }, + { + "default": "", + "name": "user_message", + "type": "string" + }, + { + "default": "", + "name": "action_required", + "type": "string" + } + ], + "name": "warmbly.events.EmailSendError", + "type": "record" + } + ] + }, + { + "default": "", + "name": "legacy_error", + "type": "string" + } + ], + "name": "warmbly.events.SendEmailResult", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "worker_id", + "type": "string" + }, + { + "default": 0, + "name": "observed_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + }, + { + "default": 0, + "name": "assigned_count", + "type": "int" + }, + { + "default": 0, + "name": "imap_idle_count", + "type": "int" + }, + { + "default": 0, + "name": "memory_mb", + "type": "int" + }, + { + "default": 0, + "name": "goroutine_count", + "type": "int" + }, + { + "default": 0, + "name": "sends_attempted", + "type": "int" + }, + { + "default": 0, + "name": "sends_succeeded", + "type": "int" + }, + { + "default": 0, + "name": "bounces_hard", + "type": "int" + }, + { + "default": 0, + "name": "bounces_soft", + "type": "int" + }, + { + "default": 0, + "name": "complaints", + "type": "int" + }, + { + "default": 0, + "name": "auth_errors", + "type": "int" + }, + { + "default": 0, + "name": "rate_limit_errors", + "type": "int" + }, + { + "default": 0, + "name": "smtp_latency_p50_ms", + "type": "int" + }, + { + "default": 0, + "name": "smtp_latency_p99_ms", + "type": "int" + } + ], + "name": "warmbly.events.WorkerHealthSample", + "type": "record" + } + ] + } + ], + "name": "warmbly.events.JobEvent", + "type": "record" +} diff --git a/internal/app/eventschemas/testdata/warmup-events.avsc b/internal/app/eventschemas/testdata/warmup-events.avsc new file mode 100644 index 000000000..6e86fd745 --- /dev/null +++ b/internal/app/eventschemas/testdata/warmup-events.avsc @@ -0,0 +1,44 @@ +{ + "fields": [ + { + "default": "", + "name": "event_type", + "type": "string" + }, + { + "default": "", + "name": "task_id", + "type": "string" + }, + { + "default": "", + "name": "sender_account_id", + "type": "string" + }, + { + "default": "", + "name": "target_account_id", + "type": "string" + }, + { + "default": "", + "name": "message_id", + "type": "string" + }, + { + "default": false, + "name": "is_reply", + "type": "boolean" + }, + { + "default": 0, + "name": "sent_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + } + ], + "name": "warmbly.events.WarmupEmailSentEvent", + "type": "record" +} diff --git a/internal/app/eventschemas/testdata/worker-commands.avsc b/internal/app/eventschemas/testdata/worker-commands.avsc new file mode 100644 index 000000000..59ae48a14 --- /dev/null +++ b/internal/app/eventschemas/testdata/worker-commands.avsc @@ -0,0 +1,839 @@ +{ + "fields": [ + { + "name": "type", + "type": "string" + }, + { + "name": "body", + "type": [ + { + "fields": [ + { + "default": "", + "name": "id", + "type": "string" + }, + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": null, + "name": "organization_id", + "type": [ + "null", + "string" + ] + }, + { + "default": false, + "name": "imap_sync", + "type": "boolean" + }, + { + "default": null, + "name": "save_to_sent", + "type": [ + "null", + "boolean" + ] + }, + { + "default": "", + "name": "email", + "type": "string" + }, + { + "default": "", + "name": "first_name", + "type": "string" + }, + { + "default": "", + "name": "last_name", + "type": "string" + }, + { + "default": "", + "name": "type", + "type": "string" + }, + { + "default": null, + "name": "google", + "type": [ + "null", + { + "fields": [ + { + "default": "\u0000\u0000\u0000\u0000\u0000\u0000\u0000\u0000", + "name": "last_history_id", + "type": { + "name": "warmbly.events.uint64", + "size": 8, + "type": "fixed" + } + }, + { + "default": null, + "name": "token", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "AccessToken", + "type": "string" + }, + { + "default": "", + "name": "TokenType", + "type": "string" + }, + { + "default": "", + "name": "RefreshToken", + "type": "string" + }, + { + "default": 0, + "name": "Expiry", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + }, + { + "default": 0, + "name": "ExpiresIn", + "type": "long" + } + ], + "name": "warmbly.events.Token", + "type": "record" + } + ] + } + ], + "name": "warmbly.events.AddWorkerEmailGoogleData", + "type": "record" + } + ] + }, + { + "default": null, + "name": "smtp_imap", + "type": [ + "null", + { + "fields": [ + { + "default": [], + "name": "mailboxes", + "type": { + "items": { + "fields": [ + { + "default": "", + "name": "name", + "type": "string" + }, + { + "default": [], + "name": "attributes", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": 0, + "name": "uid_validity", + "type": "long" + }, + { + "default": "\u0000\u0000\u0000\u0000\u0000\u0000\u0000\u0000", + "name": "highestmodseq", + "type": "warmbly.events.uint64" + }, + { + "default": 0, + "name": "uid_next", + "type": "long" + }, + { + "default": "", + "name": "delim", + "type": "string" + }, + { + "default": 0, + "name": "messages", + "type": "long" + }, + { + "default": 0, + "name": "updated_at", + "type": { + "logicalType": "timestamp-millis", + "type": "long" + } + } + ], + "name": "warmbly.events.Mailbox", + "type": "record" + }, + "type": "array" + } + }, + { + "default": null, + "name": "token", + "type": [ + "null", + "warmbly.events.Token" + ] + }, + { + "default": null, + "name": "credentials", + "type": [ + "null", + { + "fields": [ + { + "default": null, + "name": "smtp", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "username", + "type": "string" + }, + { + "default": "", + "name": "password", + "type": "string" + }, + { + "default": "", + "name": "host", + "type": "string" + }, + { + "default": 0, + "name": "port", + "type": "int" + }, + { + "default": "", + "name": "security", + "type": "string" + } + ], + "name": "warmbly.events.Service", + "type": "record" + } + ] + }, + { + "default": null, + "name": "imap", + "type": [ + "null", + "warmbly.events.Service" + ] + } + ], + "name": "warmbly.events.SmtpImap", + "type": "record" + } + ] + } + ], + "name": "warmbly.events.AddWorkerEmailSmtpImapData", + "type": "record" + } + ] + }, + { + "default": null, + "name": "graph", + "type": [ + "null", + { + "fields": [ + { + "default": null, + "name": "token", + "type": [ + "null", + "warmbly.events.Token" + ] + }, + { + "default": {}, + "name": "delta_links", + "type": { + "type": "map", + "values": "string" + } + }, + { + "default": "", + "name": "user", + "type": "string" + } + ], + "name": "warmbly.events.AddWorkerEmailGraphData", + "type": "record" + } + ] + }, + { + "default": null, + "name": "sync", + "type": [ + "null", + { + "fields": [ + { + "default": { + "backfill_days": 0, + "backfill_messages": 0, + "daily_messages": 0, + "org_daily_messages": 0, + "skip_folders": [] + }, + "name": "policy", + "type": { + "fields": [ + { + "default": 0, + "name": "backfill_days", + "type": "int" + }, + { + "default": 0, + "name": "backfill_messages", + "type": "int" + }, + { + "default": 0, + "name": "daily_messages", + "type": "int" + }, + { + "default": 0, + "name": "org_daily_messages", + "type": "int" + }, + { + "default": [], + "name": "skip_folders", + "type": { + "items": "string", + "type": "array" + } + } + ], + "name": "warmbly.events.SyncPolicy", + "type": "record" + } + }, + { + "default": null, + "name": "state", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "backfill_status", + "type": "string" + }, + { + "default": { + "folders": {}, + "page_token": "" + }, + "name": "backfill_cursor", + "type": { + "fields": [ + { + "default": "", + "name": "page_token", + "type": "string" + }, + { + "default": {}, + "name": "folders", + "type": { + "type": "map", + "values": { + "fields": [ + { + "default": "", + "name": "next", + "type": "string" + }, + { + "default": 0, + "name": "uid", + "type": "long" + }, + { + "default": false, + "name": "done", + "type": "boolean" + } + ], + "name": "warmbly.events.SyncFolderCursor", + "type": "record" + } + } + } + ], + "name": "warmbly.events.SyncCursor", + "type": "record" + } + }, + { + "default": 0, + "name": "backfill_synced", + "type": "int" + }, + { + "default": null, + "name": "backfill_since", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": null, + "name": "backfill_started_at", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": null, + "name": "backfill_completed_at", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": null, + "name": "throttled_until", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + }, + { + "default": "", + "name": "throttle_reason", + "type": "string" + }, + { + "default": 0, + "name": "deferred", + "type": "int" + }, + { + "default": 0, + "name": "folders_skipped_cap", + "type": "int" + }, + { + "default": 0, + "name": "folders_skipped_conflict", + "type": "int" + }, + { + "default": null, + "name": "last_synced_at", + "type": [ + "null", + { + "logicalType": "timestamp-millis", + "type": "long" + } + ] + } + ], + "name": "warmbly.events.SyncState", + "type": "record" + } + ] + } + ], + "name": "warmbly.events.AddWorkerEmailSyncData", + "type": "record" + } + ] + }, + { + "default": false, + "name": "brokered", + "type": "boolean" + } + ], + "name": "warmbly.events.AddWorkerEmail", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "org_id", + "type": "string" + }, + { + "default": "", + "name": "process_id", + "type": "string" + }, + { + "default": null, + "name": "credentials", + "type": [ + "null", + "warmbly.events.SmtpImap" + ] + } + ], + "name": "warmbly.events.EventWorkerEmailValidation", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "process_id", + "type": "string" + }, + { + "default": false, + "name": "want_signature", + "type": "boolean" + }, + { + "default": "", + "name": "signature_for", + "type": "string" + } + ], + "name": "warmbly.events.EventWorkerMailboxIdentity", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": false, + "name": "seen", + "type": "boolean" + }, + { + "default": [], + "name": "messages", + "type": { + "items": { + "fields": [ + { + "default": "", + "name": "provider_id", + "type": "string" + }, + { + "default": 0, + "name": "uid", + "type": "long" + }, + { + "default": "", + "name": "folder", + "type": "string" + }, + { + "default": "", + "name": "rfc_message_id", + "type": "string" + } + ], + "name": "warmbly.events.MessageSeenRef", + "type": "record" + }, + "type": "array" + } + } + ], + "name": "warmbly.events.MessageSeenAction", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + } + ], + "name": "warmbly.events.RemoveWorkerEmail", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "task_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "org_id", + "type": "string" + }, + { + "default": [], + "name": "to", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "cc", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": [], + "name": "bcc", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": "", + "name": "subject", + "type": "string" + }, + { + "default": "", + "name": "body_s3_key", + "type": "string" + }, + { + "default": "", + "name": "message_id", + "type": "string" + }, + { + "default": "", + "name": "in_reply_to", + "type": "string" + }, + { + "default": null, + "name": "parent", + "type": [ + "null", + { + "fields": [ + { + "default": "", + "name": "id", + "type": "string" + }, + { + "default": "", + "name": "message_id", + "type": "string" + }, + { + "default": "", + "name": "thread_id", + "type": "string" + } + ], + "name": "warmbly.events.EmailParent", + "type": "record" + } + ] + }, + { + "default": false, + "name": "is_warmup", + "type": "boolean" + }, + { + "default": null, + "name": "tracking_info", + "type": [ + "null", + { + "fields": [ + { + "default": false, + "name": "open_tracking", + "type": "boolean" + }, + { + "default": false, + "name": "link_tracking", + "type": "boolean" + }, + { + "default": "", + "name": "tracking_domain", + "type": "string" + } + ], + "name": "warmbly.events.TrackingInfo", + "type": "record" + } + ] + }, + { + "default": "", + "name": "warmup_token", + "type": "string" + }, + { + "default": "", + "name": "unsubscribe_url", + "type": "string" + } + ], + "name": "warmbly.events.SendEmail", + "type": "record" + }, + { + "fields": [ + { + "default": "", + "name": "user_id", + "type": "string" + }, + { + "default": "", + "name": "email_id", + "type": "string" + }, + { + "default": "", + "name": "gmail_id", + "type": "string" + }, + { + "default": 0, + "name": "uid", + "type": "long" + }, + { + "default": 0, + "name": "mailbox_uid_validity", + "type": "long" + }, + { + "default": "", + "name": "mailbox_folder", + "type": "string" + }, + { + "default": "", + "name": "rfc_message_id", + "type": "string" + }, + { + "default": [], + "name": "actions", + "type": { + "items": "string", + "type": "array" + } + }, + { + "default": "", + "name": "placement", + "type": "string" + }, + { + "default": "", + "name": "target_folder", + "type": "string" + }, + { + "default": "", + "name": "internal_id", + "type": "string" + }, + { + "default": false, + "name": "recheck", + "type": "boolean" + }, + { + "default": 0, + "name": "delay_seconds", + "type": "int" + } + ], + "name": "warmbly.events.WarmupEmailAction", + "type": "record" + } + ] + } + ], + "name": "warmbly.events.WorkerEvent", + "type": "record" +} diff --git a/internal/app/releases/service.go b/internal/app/releases/service.go index c6e81a5dd..c79c30de7 100644 --- a/internal/app/releases/service.go +++ b/internal/app/releases/service.go @@ -33,6 +33,8 @@ type Config struct { WorkerImageRepo string // "ghcr.io/warmbly/warmbly/worker" GithubToken string // optional, raises API rate limit HTTPClient *http.Client + // SchemaGate refuses a tag whose bus schemas are not registered; nil allows any. + SchemaGate func(ctx context.Context, tag string) error } type Service struct { @@ -139,6 +141,14 @@ func (s *Service) CheckGitHub(ctx context.Context) (*models.FleetReleaseState, e return current, nil } + if s.cfg.SchemaGate != nil { + if err := s.cfg.SchemaGate(ctx, head.TagName); err != nil { + s.recordError(err.Error()) + log.Printf("releases: fleet held at %s: %v", current.Tag, err) + return current, nil + } + } + next := &models.FleetReleaseState{ Channel: current.Channel, Tag: head.TagName, diff --git a/internal/app/releases/service_test.go b/internal/app/releases/service_test.go new file mode 100644 index 000000000..541c0fc2b --- /dev/null +++ b/internal/app/releases/service_test.go @@ -0,0 +1,60 @@ +package releases + +import ( + "context" + "errors" + "io" + "net/http" + "strings" + "testing" + + "github.com/warmbly/warmbly/internal/models" +) + +type memSettings struct{ release *models.FleetReleaseState } + +func (m *memSettings) GetRelease(context.Context) (*models.FleetReleaseState, error) { + return m.release, nil +} + +func (m *memSettings) SetRelease(_ context.Context, s *models.FleetReleaseState) error { + m.release = s + return nil +} + +func (m *memSettings) GetJoinTokenHash(context.Context) (string, error) { return "", nil } +func (m *memSettings) SetJoinTokenHash(context.Context, string) error { return nil } + +type githubStub struct{} + +func (githubStub) RoundTrip(*http.Request) (*http.Response, error) { + body := `[{"tag_name":"v2.0.0","published_at":"2026-09-29T12:00:00Z"}]` + return &http.Response{StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(body)), Header: http.Header{}}, nil +} + +func TestCheckGitHubHoldsTheFleetWhenTheSchemaGateRefuses(t *testing.T) { + for _, tc := range []struct { + name string + gate func(context.Context, string) error + want string + }{ + {"refused", func(context.Context, string) error { return errors.New("incompatible") }, "v1.0.0"}, + {"accepted", func(context.Context, string) error { return nil }, "v2.0.0"}, + } { + t.Run(tc.name, func(t *testing.T) { + settings := &memSettings{release: &models.FleetReleaseState{Channel: models.FleetChannelStable, Tag: "v1.0.0"}} + svc := New(Config{ + Enabled: true, + GithubRepo: "warmbly/warmbly", + HTTPClient: &http.Client{Transport: githubStub{}}, + SchemaGate: tc.gate, + }, settings) + if _, err := svc.CheckGitHub(context.Background()); err != nil { + t.Fatal(err) + } + if settings.release.Tag != tc.want { + t.Fatalf("fleet target %s, want %s", settings.release.Tag, tc.want) + } + }) + } +} diff --git a/internal/infrastructure/codec/avro.go b/internal/infrastructure/codec/avro.go index b2f77bf7b..4a23f41d0 100644 --- a/internal/infrastructure/codec/avro.go +++ b/internal/infrastructure/codec/avro.go @@ -11,10 +11,13 @@ import ( "encoding/binary" "errors" "fmt" + "strings" "sync" "github.com/confluentinc/confluent-kafka-go/v2/schemaregistry" + "github.com/confluentinc/confluent-kafka-go/v2/schemaregistry/rest" "github.com/hamba/avro/v2" + "github.com/rs/zerolog/log" "github.com/warmbly/warmbly/internal/models" ) @@ -92,6 +95,37 @@ func (c *AvroCodec) Serialize(_ context.Context, topic string, value any) ([]byt return append(out, body...), nil } +// RegisterSchemas satisfies SchemaRegistrar. +func (c *AvroCodec) RegisterSchemas(_ context.Context, schemas map[string]avro.Schema) error { + var all []string + var errs []error + for topic, schema := range schemas { + topics := []string{topic} + if prefix, ok := strings.CutSuffix(topic, "*"); ok { + if all == nil { + var err error + if all, err = c.client.GetAllSubjects(); err != nil { + return fmt.Errorf("codec: list subjects: %w", err) + } + } + topics = topics[:0] + for _, subject := range all { + if t, ok := strings.CutSuffix(subject, "-value"); ok && strings.HasPrefix(t, prefix) { + topics = append(topics, t) + } + } + } + for _, t := range topics { + if _, err := c.register(subjectFor(t), schema); err != nil { + errs = append(errs, err) + } + } + } + return errors.Join(errs...) +} + +var _ SchemaRegistrar = (*AvroCodec)(nil) + // Deserialize decodes with the schema the payload names, which is the writer's // rather than whatever this process happens to hold. That is the whole point of // carrying the id: a consumer reads what was actually written. @@ -132,6 +166,9 @@ func (c *AvroCodec) register(subject string, schema avro.Schema) (int, error) { return id, nil } id, err = c.client.Register(subject, schemaregistry.SchemaInfo{Schema: string(doc)}, true) + if isIncompatible(err) && c.pinBackward(subject) { + id, err = c.client.Register(subject, schemaregistry.SchemaInfo{Schema: string(doc)}, true) + } if err != nil { return 0, fmt.Errorf("codec: register %s: %w", subject, err) } @@ -141,6 +178,36 @@ func (c *AvroCodec) register(subject string, schema avro.Schema) (int, error) { return id, nil } +// isIncompatible is the registry refusing a schema against an earlier version. +func isIncompatible(err error) bool { + var rerr *rest.Error + return errors.As(err, &rerr) && rerr.Code == 409 +} + +// pinBackward sets the subject to BACKWARD when its effective level refuses a +// new union branch, since every new event type adds one to the envelope. A +// level that already admits it is left alone, and so is the refusal. +func (c *AvroCodec) pinBackward(subject string) bool { + level, err := c.client.GetCompatibility(subject) + if err != nil { + if level, err = c.client.GetDefaultCompatibility(); err != nil { + return false + } + } + switch level { + case schemaregistry.Forward, schemaregistry.ForwardTransitive, + schemaregistry.Full, schemaregistry.FullTransitive: + default: + return false + } + if _, err := c.client.UpdateCompatibility(subject, schemaregistry.Backward); err != nil { + return false + } + log.Warn().Str("subject", subject).Str("was", level.String()). + Msg("schema registry: subject set to BACKWARD so a new event type can be registered") + return true +} + func (c *AvroCodec) schemaByID(subject string, id int) (avro.Schema, error) { c.mu.RLock() schema, ok := c.schemas[id] diff --git a/internal/infrastructure/codec/avro_test.go b/internal/infrastructure/codec/avro_test.go index dcecbc087..1a3407edf 100644 --- a/internal/infrastructure/codec/avro_test.go +++ b/internal/infrastructure/codec/avro_test.go @@ -4,7 +4,14 @@ package codec import ( "context" + "errors" + "slices" "testing" + + "github.com/confluentinc/confluent-kafka-go/v2/schemaregistry" + "github.com/confluentinc/confluent-kafka-go/v2/schemaregistry/rest" + "github.com/hamba/avro/v2" + "github.com/warmbly/warmbly/internal/models" ) // Full round-trip coverage for AvroCodec requires a live Confluent Schema @@ -46,3 +53,91 @@ func TestAvroCodec_NilReceiverIsSafe(t *testing.T) { t.Fatal("expected error on nil receiver") } } + +// fakeRegistry refuses a new schema while the subject's effective level is +// in the forward family, the way a registry does for a new union branch. +type fakeRegistry struct { + schemaregistry.Client + global schemaregistry.Compatibility + subject schemaregistry.Compatibility + subjects []string + registered []string +} + +func (f *fakeRegistry) GetAllSubjects() ([]string, error) { return f.subjects, nil } + +func (f *fakeRegistry) effective() schemaregistry.Compatibility { + if f.subject != 0 { + return f.subject + } + return f.global +} + +func (f *fakeRegistry) Register(subject string, _ schemaregistry.SchemaInfo, _ bool) (int, error) { + f.registered = append(f.registered, subject) + switch f.effective() { + case schemaregistry.Forward, schemaregistry.Full: + return 0, &rest.Error{Code: 409, Message: "incompatible"} + } + return 7, nil +} + +func (f *fakeRegistry) GetCompatibility(string) (schemaregistry.Compatibility, error) { + if f.subject == 0 { + return 0, &rest.Error{Code: 40408, Message: "no subject-level compatibility"} + } + return f.subject, nil +} + +func (f *fakeRegistry) GetDefaultCompatibility() (schemaregistry.Compatibility, error) { + return f.global, nil +} + +func (f *fakeRegistry) UpdateCompatibility(_ string, c schemaregistry.Compatibility) (schemaregistry.Compatibility, error) { + f.subject = c + return c, nil +} + +func TestAvroCodec_RegisterPinsForwardSubjectBackward(t *testing.T) { + for _, global := range []schemaregistry.Compatibility{schemaregistry.Forward, schemaregistry.Full} { + reg := &fakeRegistry{global: global} + c := &AvroCodec{client: reg, ids: map[string]int{}, schemas: map[int]avro.Schema{}} + ev := models.JobEvent{Type: models.JobEventTypeEmailSent, Body: models.SendEmailResult{}} + if _, err := c.Serialize(context.Background(), "jobs.worker-events", ev); err != nil { + t.Fatalf("global %s: %v", global.String(), err) + } + if reg.subject != schemaregistry.Backward { + t.Fatalf("global %s: subject left at %v", global.String(), reg.subject) + } + } +} + +func TestAvroCodec_RegisterLeavesOtherRefusalsAlone(t *testing.T) { + reg := &fakeRegistry{global: schemaregistry.Forward, subject: schemaregistry.BackwardTransitive} + c := &AvroCodec{client: reg, ids: map[string]int{}, schemas: map[int]avro.Schema{}} + if !isIncompatible(&rest.Error{Code: 409}) || isIncompatible(errors.New("x")) { + t.Fatal("isIncompatible misreads the refusal") + } + if c.pinBackward("s") || reg.subject != schemaregistry.BackwardTransitive { + t.Fatal("a level that admits a new branch was changed") + } +} + +func TestAvroCodec_RegisterSchemasExpandsPerNodeTopics(t *testing.T) { + reg := &fakeRegistry{global: schemaregistry.Backward, subjects: []string{ + "w.a-value", "w.b-value", "warmup-events-value", "jobs.worker-events-value", + }} + c := &AvroCodec{client: reg, ids: map[string]int{}, schemas: map[int]avro.Schema{}} + err := c.RegisterSchemas(context.Background(), map[string]avro.Schema{ + "w.*": models.WorkerEvent{}.Schema(), + "jobs.worker-events": models.JobEvent{}.Schema(), + }) + if err != nil { + t.Fatal(err) + } + slices.Sort(reg.registered) + want := []string{"jobs.worker-events-value", "w.a-value", "w.b-value"} + if !slices.Equal(reg.registered, want) { + t.Fatalf("registered %v, want %v", reg.registered, want) + } +} diff --git a/internal/infrastructure/codec/codec.go b/internal/infrastructure/codec/codec.go index f892b0106..77d3d7d60 100644 --- a/internal/infrastructure/codec/codec.go +++ b/internal/infrastructure/codec/codec.go @@ -15,7 +15,11 @@ // in-band marker, this is an operator-level configuration choice. package codec -import "context" +import ( + "context" + + "github.com/hamba/avro/v2" +) // Codec serializes and deserializes event payloads. Implementations may // require external services (Schema Registry for Avro) or be standalone @@ -37,6 +41,14 @@ type Codec interface { Name() string } +// SchemaRegistrar is a codec that resolves schemas against a registry, so a +// release can register what it publishes before any node publishes it. +type SchemaRegistrar interface { + // RegisterSchemas registers each schema under its topic's subject. A topic + // ending in ".*" names every topic already registered under that prefix. + RegisterSchemas(ctx context.Context, schemas map[string]avro.Schema) error +} + // Compile-time interface check. The AvroCodec check lives in avro.go (behind // the `kafka` build tag) since that type isn't compiled into the default build. var _ Codec = (*JSONCodec)(nil) diff --git a/internal/models/unibox.go b/internal/models/unibox.go index 06445c031..bb1793f72 100644 --- a/internal/models/unibox.go +++ b/internal/models/unibox.go @@ -2,6 +2,7 @@ package models import ( "slices" + "strings" "time" "github.com/google/uuid" @@ -160,6 +161,24 @@ type EmailMessageStoreData struct { CreatedAt time.Time `json:"created_at" avro:"created_at"` } +// ValidText makes every text field storable: Postgres refuses invalid UTF-8 +// and NUL, and an unlabelled 8-bit header or body carries both. +func (e *EmailMessageStoreData) ValidText() { + for _, s := range []*string{&e.FolderPath, &e.Folder, &e.ProviderFolder, &e.ThreadID, &e.MessageID, + &e.GmailID, &e.ParentID, &e.Subject, &e.Snippet, &e.BodyText} { + *s = validText(*s) + } + for _, list := range [][]string{e.Flags, e.BCC, e.CC, e.FromAddr, e.InReplyTo, e.ReplyTo, e.ToAddr} { + for i := range list { + list[i] = validText(list[i]) + } + } +} + +func validText(s string) string { + return strings.ToValidUTF8(strings.ReplaceAll(s, "\x00", ""), "\uFFFD") +} + type EmailMessageStoreDataPreview struct { ID uuid.UUID `json:"id"` EmailID uuid.UUID `json:"email_id"` diff --git a/internal/models/unibox_validtext_test.go b/internal/models/unibox_validtext_test.go new file mode 100644 index 000000000..ce7979194 --- /dev/null +++ b/internal/models/unibox_validtext_test.go @@ -0,0 +1,28 @@ +package models + +import ( + "testing" + "unicode/utf8" +) + +func TestEmailMessageStoreDataValidText(t *testing.T) { + m := &EmailMessageStoreData{ + Subject: "caf\xe9 \xe2\xa2\x77", + Snippet: "a\x00b", + BodyText: "ok ✓", + FromAddr: []string{"Ren\xe9 "}, + } + m.ValidText() + + for _, s := range append([]string{m.Subject, m.Snippet, m.BodyText}, m.FromAddr...) { + if !utf8.ValidString(s) { + t.Fatalf("still invalid: %q", s) + } + } + if m.Snippet != "ab" { + t.Fatalf("NUL kept: %q", m.Snippet) + } + if m.BodyText != "ok ✓" { + t.Fatalf("valid text changed: %q", m.BodyText) + } +} diff --git a/site/src/layouts/Layout.astro b/site/src/layouts/Layout.astro index 2eb0c59c2..0a08ddd46 100644 --- a/site/src/layouts/Layout.astro +++ b/site/src/layouts/Layout.astro @@ -217,21 +217,27 @@ const websiteJsonLd = { 'ResizeObserver loop completed with undelivered notifications.', 'ResizeObserver loop limit exceeded', ]; + // A link scanner's embedded browser (Outlook Safe Links and the like) + // rejects with this while it walks the page; it is not our code. + var scanner = 'Object Not Found Matching Id:'; + var isNoise = function (v) { + return typeof v === 'string' && (noise.indexOf(v.trim()) !== -1 || v.indexOf(scanner) !== -1); + }; var exceptions = event.properties && event.properties.$exception_list; if (Array.isArray(exceptions)) { for (var i = 0; i < exceptions.length; i++) { var value = exceptions[i] && (exceptions[i].value || exceptions[i].$exception_value); - if (typeof value === 'string' && noise.indexOf(value.trim()) !== -1) return null; + if (isNoise(value)) return null; } } var message = event.properties && event.properties.$exception_message; - if (typeof message === 'string' && noise.indexOf(message.trim()) !== -1) return null; + if (isNoise(message)) return null; // Keep accepting flattened payloads while cached SDK chunks are // still in browsers during a rolling release. var values = event.properties && event.properties.$exception_values; if (!Array.isArray(values)) return event; for (var i = 0; i < values.length; i++) { - if (typeof values[i] === 'string' && noise.indexOf(values[i].trim()) !== -1) return null; + if (isNoise(values[i])) return null; } return event; }, diff --git a/web/src/lib/passkey.ts b/web/src/lib/passkey.ts index 5d37d3859..4ce45b207 100644 --- a/web/src/lib/passkey.ts +++ b/web/src/lib/passkey.ts @@ -80,10 +80,15 @@ function mapError(e: unknown): Error { return new PasskeyCancelled("aborted"); case "NotAllowedError": return new PasskeyCancelled("not-allowed"); + // Another WebAuthn request is still pending, so this one never started. + case "InvalidStateError": + return new PasskeyCancelled("aborted"); default: return new Error(e.message || "Your device couldn't complete the passkey request."); } } + // SimpleWebAuthn's words for a ceremony the browser ended with no credential. + if (e instanceof Error && e.message === "Authentication was not completed") return new PasskeyCancelled("aborted"); if (e instanceof Error) return e; return new Error("Something went wrong with the passkey request."); } @@ -125,7 +130,7 @@ async function startExplicitAuthentication( }) as PublicKeyCredential | null; if (!credential) { - throw new Error("Authentication was not completed"); + throw new PasskeyCancelled("aborted"); } const response = credential.response as AuthenticatorAssertionResponse;