entitysync

package
v0.16.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 29, 2026 License: Apache-2.0 Imports: 16 Imported by: 0

Documentation

Overview

Package entitysync replicates schema-authorized runtime entities to Miren Cloud over the negotiated uplink.

Index

Constants

View Source
const (
	Version1 uint = 1

	TypeSnapshotBegin    = "entity.snapshot.begin"
	TypeSnapshotBatch    = "entity.snapshot.batch"
	TypeSnapshotComplete = "entity.snapshot.complete"
	TypeChangeBatch      = "entity.change.batch"
	TypeAck              = "entity.ack"
)

Variables

This section is empty.

Functions

This section is empty.

Types

type Ack

type Ack struct {
	MessageID string `json:"message_id"`
	Cursor    int64  `json:"cursor"`
	Error     string `json:"error,omitempty"`
}

type Change

type Change struct {
	Op       ChangeOp       `json:"op"`
	Revision int64          `json:"revision"`
	EntityID entity.Id      `json:"entity_id"`
	Kind     entity.Id      `json:"kind"`
	Entity   *entity.Entity `json:"entity,omitempty"`
}

type ChangeBatch

type ChangeBatch struct {
	MessageID          string   `json:"message_id"`
	FromRevision       int64    `json:"from_revision"`
	ToRevision         int64    `json:"to_revision"`
	ExportSchemaDigest string   `json:"export_schema_digest"`
	Changes            []Change `json:"changes"`
	SourceEpoch        string   `json:"source_epoch"`
}

type ChangeOp

type ChangeOp string
const (
	ChangePut    ChangeOp = "put"
	ChangeDelete ChangeOp = "delete"
)

type Config

type Config struct {
	ExportSchema              string `json:"export_schema"`
	Cursor                    int64  `json:"cursor"`
	SnapshotRequired          bool   `json:"snapshot_required,omitempty"`
	SourceEpoch               string `json:"source_epoch"`
	ResnapshotAfterSeconds    int64  `json:"resnapshot_after_seconds,omitempty"`
	ResnapshotIntervalSeconds int64  `json:"resnapshot_interval_seconds,omitempty"`
}

type Diagnostics

type Diagnostics struct {
	// contains filtered or unexported fields
}

Diagnostics collects transient state for the local debug interface.

func NewDiagnostics

func NewDiagnostics(schemaDigest string) *Diagnostics

func (*Diagnostics) LandedRevision added in v0.16.0

func (d *Diagnostics) LandedRevision() (int64, bool)

LandedRevision reports the highest store revision cloud has durably landed, and whether entity export applies to this cluster at all. An unregistered cluster reports false, so callers gating local deletion on cloud custody can proceed freely; a registered one reports true with zero until this process confirms progress.

Only a positive "disabled" turns export off. Before the uplink has been started the state is not yet known, and a cluster in that window is treated as exporting so a caller holds rather than deletes. An unregistered cluster pays for that with at most one deferred sweep before SetDisabled runs; a registered cluster whose uplink is down or still connecting would otherwise be told cloud does not apply, and prune history cloud never received.

func (d *Diagnostics) ObserveUplink(status uplink.Status)

func (*Diagnostics) SetDisabled

func (d *Diagnostics) SetDisabled(reason string)

func (*Diagnostics) SetPreparation

func (d *Diagnostics) SetPreparation(state, detail string)

func (*Diagnostics) SetPreparationFailure

func (d *Diagnostics) SetPreparationFailure(state, message string)

func (*Diagnostics) SnapshotStatus

func (d *Diagnostics) SnapshotStatus() Status

type Event

type Event struct {
	At        time.Time
	Stage     string
	Message   string
	MessageID string
	Cursor    int64
}

Event records the most recent acknowledgement or failure.

type Exporter

type Exporter struct {
	// contains filtered or unexported fields
}

func NewExporter

func NewExporter(log *slog.Logger, store entity.Store, contract *entityexport.Contract, options ...Option) *Exporter

func (*Exporter) Register

func (t *Exporter) Register(_ context.Context, link Link) error
type Link interface {
	OfferCapabilityFunc(string, []uint, uplink.CapabilityOfferFunc)
	OnSession(func(context.Context, uplink.Session))
	Handle(string, uplink.MessageHandler)
	SendMessageBlocking(context.Context, string, any) error
}

type Offer

type Offer struct {
	ExportSchemas []string `json:"export_schemas"`
	SourceEpoch   string   `json:"source_epoch"`
}

type Option

type Option func(*Exporter)

func WithDiagnostics

func WithDiagnostics(diagnostics *Diagnostics) Option

WithDiagnostics publishes exporter progress to the local debug interface.

func WithStartGate

func WithStartGate(ready <-chan struct{}) Option

WithStartGate delays entity reads and transmission until source preparation has completed. Capability negotiation still happens immediately so other capabilities on the shared uplink are never gated on entity migration.

type SnapshotBatch

type SnapshotBatch struct {
	SnapshotID         string           `json:"snapshot_id"`
	ExportSchemaDigest string           `json:"export_schema_digest"`
	Entities           []*entity.Entity `json:"entities"`
}

type SnapshotBegin

type SnapshotBegin struct {
	SnapshotID         string `json:"snapshot_id"`
	SourceHead         int64  `json:"source_head"`
	ExportSchemaDigest string `json:"export_schema_digest"`
	SourceEpoch        string `json:"source_epoch"`
}

type SnapshotComplete

type SnapshotComplete struct {
	MessageID    string           `json:"message_id"`
	SnapshotID   string           `json:"snapshot_id"`
	SourceHead   int64            `json:"source_head"`
	CountsByKind map[string]int64 `json:"counts_by_kind"`
}

type SnapshotProgress

type SnapshotProgress struct {
	ID           string
	HeadRevision int64
	NextRevision int64
	PagesSent    int64
	EntitiesSent int64
	CountsByKind map[string]int64
}

SnapshotProgress describes the snapshot currently being transmitted.

type Status

type Status struct {
	UplinkState       string
	SessionID         string
	HandshakeVersion  uint
	CapabilityState   string
	CapabilityVersion uint
	PreparationState  string
	PreparationDetail string
	SourceEpoch       string
	SchemaDigest      string
	Mode              string
	WaitReason        string
	CloudCursor       int64
	// LandedRevision is the highest store revision this process has confirmed
	// cloud holds durably: a completed snapshot, a committed change batch, or
	// a session resumed at a cursor cloud reported and the local epoch
	// validated. Unlike CloudCursor it is never set from an unvalidated
	// session config, so a consumer deciding whether the runtime may forget
	// an entity can trust it. Zero means nothing is confirmed yet.
	LandedRevision     int64
	NextWatchRevision  int64
	Snapshot           *SnapshotProgress
	LastAcknowledgment *Event
	LastError          *Event
	RetryAt            time.Time
}

Status is a point-in-time view of runtime entity sync.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL