worker

package
v0.13.2 Latest Latest
Warning

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

Go to latest
Published: Sep 23, 2026 License: MIT Imports: 17 Imported by: 0

Documentation

Overview

Package worker implements an opt-in durable personal-operations worker.

The worker intentionally owns only the literal top-level "ops" metadata value. All other card metadata remains the server's data and is never reconstructed or replaced by this package.

Index

Constants

View Source
const (
	NoticeQueued    = "queued"
	NoticeDelivered = "delivered"
)
View Source
const (
	MetadataKey   = "ops"
	SchemaVersion = 1
)

Variables

View Source
var (
	ErrLostOwnership         = errors.New("worker lost its durable run ownership")
	ErrNoticeReceiptConflict = errors.New("worker notice receipt conflicts with durable evidence")
)
View Source
var (
	ErrConflict = errors.New("worker metadata conflict")
	ErrNotFound = errors.New("worker card not found")
)
View Source
var (
	ErrOutputTooLarge = errors.New("worker adapter output exceeds configured limit")
	ErrPacketTooLarge = errors.New("worker adapter packet exceeds configured limit")
	ErrAdapterInvalid = errors.New("worker adapter returned an invalid result")
)
View Source
var ErrReceiptCapacity = errors.New("worker receipt history reached its durable capacity")

Functions

This section is empty.

Types

type ActionIntent

type ActionIntent struct {
	ID      string          `json:"id"`
	Intent  json.RawMessage `json:"intent"`
	Receipt *ActionReceipt  `json:"receipt,omitempty"`
}

ActionIntent is supplied by the enrolling operator. It is durable before an adapter starts, so a runner cannot invent or expand an outward action.

type ActionReceipt

type ActionReceipt struct {
	ActionID  string `json:"action_id"`
	ReceiptID string `json:"receipt_id"`
	Status    string `json:"status,omitempty"`
}

type AdapterResult

type AdapterResult struct {
	RunID          string         `json:"run_id"`
	Status         ResultStatus   `json:"status"`
	Summary        string         `json:"summary,omitempty"`
	WakeAt         *time.Time     `json:"wake_at,omitempty"`
	EventID        string         `json:"event_id,omitempty"`
	DecisionPrompt string         `json:"decision_prompt,omitempty"`
	ReceiptID      string         `json:"receipt_id,omitempty"`
	ActionReceipt  *ActionReceipt `json:"action_receipt,omitempty"`
	ActionStatus   string         `json:"action_status,omitempty"`
	ReviewNote     string         `json:"review_note,omitempty"`
}

func (AdapterResult) Validate

func (r AdapterResult) Validate(action *ActionIntent, now time.Time) error

type CheckReport

type CheckReport struct {
	Report
	Tasks []StatusItem `json:"tasks"`
}

type Client

type Client interface {
	GetBoard(context.Context, string, bool) (json.RawMessage, error)
	GetCard(context.Context, string) (json.RawMessage, error)
	GetCardMarkdown(context.Context, string) (string, error)
	GetCardMetadata(context.Context, string) (api.CardMetadata, error)
	UpdateCardMetadata(context.Context, string, api.MetadataPatch) (json.RawMessage, error)
	CreateCard(context.Context, string, string, string, string) (json.RawMessage, error)
	AddCommentOnce(context.Context, string, string) (json.RawMessage, error)
}

Client is intentionally smaller than api.Client so worker tests can use a real HTTP-shaped fake without exposing unrelated board operations.

type Clock

type Clock interface{ Now() time.Time }

type Config

type Config struct {
	BoardID       string
	WorkerID      string
	PollInterval  time.Duration
	LeaseDuration time.Duration
	RunTimeout    time.Duration
	MaxConcurrent int
	ArtifactDir   string
	OutputLimit   int
	PacketLimit   int
	NoticeTimeout time.Duration
	Runner        Runner
	Notifier      Notifier
	// Observer is optional and is run only by Serve's shared board-wide pass.
	// Its proposals remain non-delegated and are never passed to Runner.
	Observer               Observer
	ObserverRegistryCardID string
	ObserverListID         string
	ObserverTimeout        time.Duration
}

func (Config) ValidateBase

func (c Config) ValidateBase() error

func (Config) ValidateExecution

func (c Config) ValidateExecution() error

func (Config) ValidateReadOnly

func (c Config) ValidateReadOnly() error

func (Config) ValidateScheduledObserver

func (c Config) ValidateScheduledObserver() error

type Decision

type Decision struct {
	Prompt string          `json:"prompt,omitempty"`
	ID     string          `json:"id,omitempty"`
	Value  json.RawMessage `json:"value,omitempty"`
}

type Delegation

type Delegation struct {
	Delegated     bool            `json:"delegated"`
	Authorization json.RawMessage `json:"authorization,omitempty"`
}

type Enrollment

type Enrollment struct {
	Goal               string
	CompletionCriteria string
	Authorization      json.RawMessage
	Sources            []SourceRef
	Action             *ActionIntent
}

type IngestReport

type IngestReport struct {
	Received   int      `json:"received"`
	Created    []string `json:"created"`
	Duplicates int      `json:"duplicates"`
}

type Journal

type Journal struct {
	ID        string `json:"id"`
	State     string `json:"state"`
	CommentID string `json:"comment_id,omitempty"`
}

type LeaseTicker

type LeaseTicker interface {
	C() <-chan time.Time
	Stop()
}

type Notice

type Notice struct {
	ID    string `json:"id"`
	State string `json:"state"`
	// QueueReceiptID proves only that a relay accepted the notice for later
	// display. It is intentionally distinct from a user-delivery receipt.
	QueueReceiptID string `json:"queue_receipt_id,omitempty"`
	ReceiptID      string `json:"receipt_id,omitempty"`
}

func (Notice) Validate

func (n Notice) Validate() error

Validate rejects malformed durable notice evidence. A queued receipt is transport acceptance only, while a delivered receipt is final proof from a synchronous adapter or a later operator acknowledgement.

type NoticePacket

type NoticePacket struct {
	Version  int    `json:"version"`
	NoticeID string `json:"notice_id"`
	CardID   string `json:"card_id"`
	RunID    string `json:"run_id"`
	// DecisionID is sent only with a waiting-user outcome, allowing a host
	// relay to route the displayed question and answer exactly.
	DecisionID string `json:"decision_id,omitempty"`
	Message    string `json:"message"`
}

type NoticeReceipt

type NoticeReceipt struct {
	NoticeID  string `json:"notice_id"`
	ReceiptID string `json:"receipt_id"`
	// An omitted state retains synchronous-adapter compatibility as delivered.
	// queued means the adapter accepted transport, not that a user saw it.
	DeliveryState string `json:"delivery_state,omitempty"`
}

func (*NoticeReceipt) ValidateFor

func (r *NoticeReceipt) ValidateFor(noticeID string) error

ValidateFor normalizes legacy synchronous receipts and checks every adapter result. Service calls this too, because custom Notifier implementations do not pass through the subprocess JSON boundary.

type Notifier

type Notifier interface {
	Deliver(context.Context, NoticePacket) (NoticeReceipt, error)
}

type Observer

type Observer interface {
	Observe(context.Context) ([]Suggestion, error)
}

type Outcome

type Outcome struct {
	RunID      string `json:"run_id"`
	Kind       string `json:"kind"`
	Summary    string `json:"summary,omitempty"`
	ReceiptID  string `json:"receipt_id,omitempty"`
	ReviewNote string `json:"review_note,omitempty"`
}

type Packet

type Packet struct {
	Version            int             `json:"version"`
	CardID             string          `json:"card_id"`
	BoardID            string          `json:"board_id"`
	RunID              string          `json:"run_id"`
	Goal               string          `json:"goal"`
	CompletionCriteria string          `json:"completion_criteria"`
	Authorization      json.RawMessage `json:"authorization"`
	CardMarkdown       string          `json:"card_markdown"`
	Sources            []SourceRef     `json:"sources,omitempty"`
	EventReceipts      []string        `json:"event_receipts,omitempty"`
	Decision           *Decision       `json:"decision,omitempty"`
	Action             *ActionIntent   `json:"action,omitempty"`
	AllowedResults     []string        `json:"allowed_results"`
}

type Record

type Record struct {
	Version            int           `json:"version"`
	Goal               string        `json:"goal"`
	CompletionCriteria string        `json:"completion_criteria"`
	Delegation         Delegation    `json:"delegation"`
	State              State         `json:"state"`
	WakeAt             *time.Time    `json:"wake_at,omitempty"`
	Sources            []SourceRef   `json:"sources,omitempty"`
	Action             *ActionIntent `json:"action,omitempty"`
	Run                *Run          `json:"run,omitempty"`
	Outcome            *Outcome      `json:"outcome,omitempty"`
	Decision           *Decision     `json:"decision,omitempty"`
	EventID            string        `json:"event_id,omitempty"`
	Journal            *Journal      `json:"journal,omitempty"`
	Notice             *Notice       `json:"notice,omitempty"`
	// Unresolved outbox entries are retained when a later task step creates a
	// new current journal/notice. They are evidence for operator reconciliation
	// and are never replayed automatically.
	UnresolvedJournals []Journal                  `json:"unresolved_journals,omitempty"`
	UnresolvedNotices  []Notice                   `json:"unresolved_notices,omitempty"`
	EventReceipts      []string                   `json:"event_receipts,omitempty"`
	DecisionReceipts   []string                   `json:"decision_receipts,omitempty"`
	SuggestionRegistry bool                       `json:"suggestion_registry,omitempty"`
	SuggestionClaims   map[string]SuggestionClaim `json:"suggestion_claims,omitempty"`
}

Record is the versioned metadata value at MetadataKey. It intentionally contains no provider connection or secret fields.

func ParseRecord

func ParseRecord(raw json.RawMessage) (Record, error)

func (*Record) AddDecisionReceipt

func (r *Record) AddDecisionReceipt(id string) error

func (*Record) AddEventReceipt

func (r *Record) AddEventReceipt(id string) error

func (Record) Eligible

func (r Record) Eligible(now time.Time) bool

func (Record) HasDecisionReceipt

func (r Record) HasDecisionReceipt(id string) bool

func (Record) HasEventReceipt

func (r Record) HasEventReceipt(id string) bool

func (*Record) RetainUnresolvedOutbox

func (r *Record) RetainUnresolvedOutbox()

func (Record) Validate

func (r Record) Validate() error

type Report

type Report struct {
	Scanned     int `json:"scanned"`
	Recognized  int `json:"recognized"`
	Due         int `json:"due"`
	Claimed     int `json:"claimed"`
	Processed   int `json:"processed"`
	NeedsReview int `json:"needs_review"`
	Ignored     int `json:"ignored"`
}

type ResultStatus

type ResultStatus string
const (
	ResultCompleted    ResultStatus = "completed"
	ResultScheduled    ResultStatus = "scheduled"
	ResultWaitingEvent ResultStatus = "waiting_event"
	ResultWaitingUser  ResultStatus = "waiting_user"
	ResultNeedsReview  ResultStatus = "needs_review"
)

type Run

type Run struct {
	ID         string `json:"id"`
	OwnerToken string `json:"owner_token"`
	// Fence binds a claim to the authorization and outward-action intent that
	// the operator approved. A metadata edit cannot retain an old claim while
	// changing either permission boundary.
	Fence        string    `json:"fence"`
	StartedAt    time.Time `json:"started_at"`
	HeartbeatAt  time.Time `json:"heartbeat_at"`
	LeaseExpires time.Time `json:"lease_expires_at"`
}

type Runner

type Runner interface {
	Run(context.Context, Packet) (AdapterResult, error)
}

type Service

type Service struct {
	Store  Store
	Config Config
	Clock  Clock
	// NewLeaseTicker is injectable so fencing/cancellation tests do not need
	// wall-clock sleeps. Production leaves it nil and uses time.NewTicker.
	NewLeaseTicker func(time.Duration) LeaseTicker
}

func (Service) AcknowledgeNotice

func (s Service) AcknowledgeNotice(ctx context.Context, cardID, noticeID, queueReceiptID, receiptID string) (bool, error)

AcknowledgeNotice records a host relay's later proof that one exact queued notice reached its conversation. It uses metadata CAS only: no Runner, Notifier, comment, or external adapter is invoked.

func (Service) Check

func (s Service) Check(ctx context.Context) (CheckReport, error)

func (Service) Decide

func (s Service) Decide(ctx context.Context, cardID, decisionID string, value json.RawMessage) (bool, error)

Decide records an explicit operator-provided decision before making the intended waiting card eligible. A runner sees that decision in its next packet but cannot write its own delegation.

func (Service) Enroll

func (s Service) Enroll(ctx context.Context, cardID string, enrollment Enrollment) error

func (Service) IngestSuggestions

func (s Service) IngestSuggestions(ctx context.Context, registryCardID, listID string, suggestions []Suggestion) (IngestReport, error)

IngestSuggestions only creates non-delegated suggested records. It has no Runner dependency by construction and cannot promote an observer proposal.

func (Service) InitSuggestionRegistry

func (s Service) InitSuggestionRegistry(ctx context.Context, cardID string) error

InitSuggestionRegistry marks an existing operator-selected card as the board's source-id reservation ledger. It is paused and non-delegated, so it cannot ever be picked up by a task runner.

func (Service) RunOnce

func (s Service) RunOnce(ctx context.Context) (Report, error)

func (Service) Serve

func (s Service) Serve(ctx context.Context) error

func (Service) Wake

func (s Service) Wake(ctx context.Context, cardID, eventID string) (bool, error)

Wake accepts a stable event receipt and only changes the target card. A duplicate receipt is a no-op even when a later polling pass occurs.

type Snapshot

type Snapshot struct {
	Card      cardView
	Metadata  api.CardMetadata
	Record    Record
	HasRecord bool
	ParseErr  error
}

type SourceRef

type SourceRef struct {
	ID        string `json:"id"`
	Reference string `json:"reference,omitempty"`
}

type State

type State string

State is a durable lifecycle state. Only ready and due scheduled records can be claimed. In particular, a stale running record is reconciled rather than treated as ready work.

const (
	StateSuggested    State = "suggested"
	StateReady        State = "ready"
	StateScheduled    State = "scheduled"
	StateRunning      State = "running"
	StateWaitingEvent State = "waiting_event"
	StateWaitingUser  State = "waiting_user"
	StateCompleted    State = "completed"
	StateCancelled    State = "cancelled"
	StatePaused       State = "paused"
	StateNeedsReview  State = "needs_review"
)

type StatusItem

type StatusItem struct {
	CardID string     `json:"card_id"`
	State  State      `json:"state"`
	WakeAt *time.Time `json:"wake_at,omitempty"`
}

type Store

type Store struct {
	Client Client
}

func (Store) CreateSuggested

func (s Store) CreateSuggested(ctx context.Context, boardID, listID string, proposal Suggestion) (string, error)

func (Store) Live

func (s Store) Live(ctx context.Context, boardID, cardID string) (Snapshot, error)

func (Store) Load

func (s Store) Load(ctx context.Context, cardID string) (Snapshot, error)

func (Store) Put

func (s Store) Put(ctx context.Context, snapshot Snapshot, record Record) (Snapshot, error)

Put changes only the worker's literal top-level metadata key. api metadata writes retain every unrelated key and enforce the provided revision.

func (Store) Scan

func (s Store) Scan(ctx context.Context, boardID string) ([]Snapshot, error)

type SubprocessNotifier

type SubprocessNotifier struct {
	Argv        []string
	OutputLimit int
	Env         []string
}

func (SubprocessNotifier) Deliver

type SubprocessObserver

type SubprocessObserver struct {
	Argv        []string
	OutputLimit int
	MaxEvents   int
	Env         []string
}

func (SubprocessObserver) Observe

func (o SubprocessObserver) Observe(ctx context.Context) ([]Suggestion, error)

type SubprocessRunner

type SubprocessRunner struct {
	Argv        []string
	ArtifactDir string
	OutputLimit int
	PacketLimit int
	// Env is deliberately explicit. Ambient process credentials are never
	// inherited by adapters; trusted bridge configuration belongs in files or
	// a socket owned by that bridge, not this worker's environment.
	Env []string
}

func (SubprocessRunner) Run

func (r SubprocessRunner) Run(ctx context.Context, packet Packet) (AdapterResult, error)

type Suggestion

type Suggestion struct {
	SourceID  string `json:"source_id"`
	Reference string `json:"reference,omitempty"`
	Title     string `json:"title"`
	Proposal  string `json:"proposal"`
}

type SuggestionClaim

type SuggestionClaim struct {
	State  string `json:"state"`
	CardID string `json:"card_id,omitempty"`
}

SuggestionClaim is an idempotency reservation stored on an explicitly configured registry card. It is deliberately separate from a suggestion card: a create response can be ambiguous, in which case the reservation is held for review instead of risking a second source card.

Jump to

Keyboard shortcuts

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