mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-03 16:02:02 +00:00
* feat: make both bus envelopes Avro-encodable by deriving each one's schema from a declared registry of body types, with a union branch per body and our own struct walk that skips unexported fields and honours avro:"-" before descending, so the schema describes exactly what encoding/json already puts on the wire, and narrow the two sync cursors on the wire DTOs to int64 because Avro has no unsigned 64-bit type * feat: stop the instance health check, the config registry and the docs all claiming Avro cannot serialize a worker envelope, which stopped being true once the envelopes carried a declared union, and check the one thing that is still a real misconfiguration instead: avro selected with no SCHEMA_REGISTRY_URL to resolve against * feat: emit a reference the second time a record appears in an envelope schema instead of defining it again, because Avro names a record once and a document that defines warmbly.events.Token three times is rejected outright, and keep the Schema Registry round-trip as a skip-by-default test since only a registry judges the document rather than the objects it was built from * feat: carry uint64 as Avro fixed(8) rather than long, which lets the sync cursors keep their unsigned type instead of being narrowed, name every event field after its json tag so the schema and the JSON wire agree, and populate every field in the round-trip test because zero values are why a uint64 mapped to long passed in the first place * feat: frame Avro in Confluent's wire format and encode through hamba's default API instead of going through avrov2, whose private avro.API holds a type resolver avro.Register cannot reach, so a union body failed there with unable to resolve type while encoding cleanly against the same schema, and keep the registry round-trip as a skip-by-default test
119 lines
5.5 KiB
Go
119 lines
5.5 KiB
Go
package models
|
|
|
|
import (
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
)
|
|
|
|
// SyncPolicy is the fair-use budget a mailbox syncs under. The control plane
|
|
// resolves it when the mailbox is loaded onto a worker (instance settings,
|
|
// see internal/app/instancesettings) and ships it inside ADD_EMAIL, so the
|
|
// worker never needs to read configuration of its own.
|
|
type SyncPolicy struct {
|
|
// BackfillDays is how far back the initial import reaches.
|
|
BackfillDays int `json:"backfill_days" avro:"backfill_days"`
|
|
// BackfillMessages caps how many messages the initial import stores.
|
|
BackfillMessages int `json:"backfill_messages" avro:"backfill_messages"`
|
|
// DailyMessages caps new (live) messages stored per UTC day. Replies to
|
|
// the mailbox's own sends have a separate budget of the same size.
|
|
DailyMessages int `json:"daily_messages" avro:"daily_messages"`
|
|
// OrgDailyMessages caps new plus backfilled messages stored across the
|
|
// whole organization per UTC day.
|
|
OrgDailyMessages int `json:"org_daily_messages" avro:"org_daily_messages"`
|
|
}
|
|
|
|
// SyncBackfillStatus is where the initial import stands.
|
|
type SyncBackfillStatus string
|
|
|
|
const (
|
|
// SyncBackfillPending: the mailbox is loaded but the import has not run yet.
|
|
SyncBackfillPending SyncBackfillStatus = "pending"
|
|
// SyncBackfillRunning: the import is walking history and may be paced.
|
|
SyncBackfillRunning SyncBackfillStatus = "running"
|
|
// SyncBackfillComplete: the window is exhausted or the cap was reached.
|
|
SyncBackfillComplete SyncBackfillStatus = "complete"
|
|
)
|
|
|
|
// SyncFolderCursor is the resumable position inside one folder of a backfill.
|
|
// Both IMAP and Graph key folders by name and IMAP walks UIDs downward, Graph
|
|
// follows @odata.nextLink; Gmail has no folders and uses
|
|
// SyncCursor.PageToken.
|
|
type SyncFolderCursor struct {
|
|
// Next is an opaque continuation (Graph nextLink).
|
|
Next string `json:"next,omitempty" avro:"next"`
|
|
// UID is the lowest UID already imported; the walk continues below it.
|
|
UID uint32 `json:"uid,omitempty" avro:"uid"`
|
|
// Done marks the folder exhausted for this window.
|
|
Done bool `json:"done,omitempty" avro:"done"`
|
|
}
|
|
|
|
// SyncCursor is the resumable position of a mailbox's backfill. It is stored
|
|
// as jsonb: read-then-execute state that is never filtered in SQL.
|
|
type SyncCursor struct {
|
|
// PageToken is Gmail's messages.list continuation.
|
|
PageToken string `json:"page_token,omitempty" avro:"page_token"`
|
|
// Folders is the per-folder position for IMAP and Graph.
|
|
Folders map[string]SyncFolderCursor `json:"folders,omitempty" avro:"folders"`
|
|
}
|
|
|
|
// SyncState is what the platform knows about a mailbox's sync: backfill
|
|
// progress, whether fair use is holding it, and when it last ran. The worker
|
|
// owns the live copy and relays every change as SYNC_STATE; the consumer
|
|
// persists it and the loader hands it back on the next (re)assignment.
|
|
type SyncState struct {
|
|
BackfillStatus SyncBackfillStatus `json:"backfill_status" avro:"backfill_status"`
|
|
BackfillCursor SyncCursor `json:"backfill_cursor" avro:"backfill_cursor"`
|
|
// BackfillSynced counts messages the import has stored so far.
|
|
BackfillSynced int `json:"backfill_synced" avro:"backfill_synced"`
|
|
// BackfillSince is the cutoff the running import uses; fixed at start so
|
|
// a later settings change does not move the goalposts mid-walk.
|
|
BackfillSince *time.Time `json:"backfill_since,omitempty" avro:"backfill_since"`
|
|
BackfillStartedAt *time.Time `json:"backfill_started_at,omitempty" avro:"backfill_started_at"`
|
|
BackfillCompletedAt *time.Time `json:"backfill_completed_at,omitempty" avro:"backfill_completed_at"`
|
|
|
|
// ThrottledUntil is set while fair use is deferring live mail; nil when
|
|
// the mailbox is within budget.
|
|
ThrottledUntil *time.Time `json:"throttled_until,omitempty" avro:"throttled_until"`
|
|
// ThrottleReason names the exhausted budget (see SyncThrottle* constants).
|
|
ThrottleReason string `json:"throttle_reason,omitempty" avro:"throttle_reason"`
|
|
// Deferred counts live messages currently waiting on budget: seen on the
|
|
// server but not yet stored. Drops back to zero once they are admitted.
|
|
Deferred int `json:"deferred" avro:"deferred"`
|
|
|
|
// FoldersSkippedCap and FoldersSkippedConflict are what the last folder
|
|
// listing could not follow: more folders than the sync covers, and
|
|
// folders whose name the server listed more than once. Carried as state
|
|
// rather than raised as an error once, so the warning goes away by itself
|
|
// when the user fixes it.
|
|
FoldersSkippedCap int `json:"folders_skipped_cap,omitempty" avro:"folders_skipped_cap"`
|
|
FoldersSkippedConflict int `json:"folders_skipped_conflict,omitempty" avro:"folders_skipped_conflict"`
|
|
|
|
LastSyncedAt *time.Time `json:"last_synced_at,omitempty" avro:"last_synced_at"`
|
|
}
|
|
|
|
// Throttle reasons, stable strings the dashboard maps to copy.
|
|
const (
|
|
SyncThrottleBurst = "burst"
|
|
SyncThrottleHourly = "hourly"
|
|
SyncThrottleDaily = "daily"
|
|
SyncThrottleOrgDaily = "org_daily"
|
|
SyncThrottlePriorityFull = "priority_daily"
|
|
)
|
|
|
|
// AddWorkerEmailSyncData seeds a mailbox's sync on the worker: the budget it
|
|
// runs under and where a previous worker left off. State is nil on first
|
|
// connect.
|
|
type AddWorkerEmailSyncData struct {
|
|
Policy SyncPolicy `json:"policy" avro:"policy"`
|
|
State *SyncState `json:"state" avro:"state"`
|
|
}
|
|
|
|
// JobEventSyncState is the worker's SYNC_STATE relay: the full current state,
|
|
// not a delta, so a lost event is repaired by the next one.
|
|
type JobEventSyncState struct {
|
|
UserID uuid.UUID `json:"user_id" avro:"user_id"`
|
|
EmailID uuid.UUID `json:"email_id" avro:"email_id"`
|
|
State SyncState `json:"state" avro:"state"`
|
|
}
|