feat: register the release's bus schemas from the booting backend and hold the fleet on its current release until the backend runs the new one and the schema registry accepts them, check every published schema against a recorded snapshot in CI (make schemas) so a change the registry would refuse fails before merge, pin a FORWARD or FULL registry subject to BACKWARD on refusal, bound advisor narration so the run still lands, store unlabelled 8-bit mail with its invalid UTF-8 replaced, treat a pending or credential-less passkey ceremony as a cancellation, and drop link-scanner rejections from site error tracking

This commit is contained in:
Matthew Meszaros
2026-09-29 19:20:07 +02:00
parent 06db0c7cf5
commit ca66b25d8d
23 changed files with 2531 additions and 11 deletions
+4
View File
@@ -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:
+6 -1
View File
@@ -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
+12
View File
@@ -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,
)
@@ -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`.
@@ -312,7 +312,7 @@ A worker's command topic is named after the node id issued when it joined, so th
</Callout>
`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.
<Callout type="warn" title="The tracking topic is read by two languages">
`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.
+4 -2
View File
@@ -20,13 +20,15 @@ Message encoding is orthogonal to transport, selected by `CODEC_PROVIDER`:
<Callout type="warn" title="Every derived field carries a default">
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.
</Callout>
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.
<Callout type="warn" title="One codec covers both topics">
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.
</Callout>
## Topics
+8 -2
View File
@@ -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 {
+1
View File
@@ -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)
+67
View File
@@ -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.<node-id>).
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")
}
@@ -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 := &registry{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.<node-id>", WorkerCommands)
}
}
+59
View File
@@ -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"
}
File diff suppressed because it is too large Load Diff
+44
View File
@@ -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"
}
+839
View File
@@ -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"
}
+10
View File
@@ -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,
+60
View File
@@ -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)
}
})
}
}
+67
View File
@@ -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]
@@ -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)
}
}
+13 -1
View File
@@ -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)
+19
View File
@@ -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"`
+28
View File
@@ -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 <rene@example.com>"},
}
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)
}
}
+9 -3
View File
@@ -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;
},
+6 -1
View File
@@ -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;