jobs

package
v1.1.0 Latest Latest
Warning

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

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

Documentation

Overview

Package jobs owns durable scheduling over the original Datly DATLY_JOBS table.

Index

Constants

View Source
const MaxInlineState = 63 * 1024

Variables

View Source
var ErrCompletionPending = errors.New("job completion requires managed transaction reconciliation")
View Source
var ErrInProgress = errors.New("job execution is already in progress")
View Source
var ErrJobFailed = errors.New("job execution failed")
View Source
var ErrNotFound = errors.New("job not found")
View Source
var ErrResultUnavailable = errors.New("completed job result is not retained; use a new match key for synchronous execution")

ErrResultUnavailable distinguishes durable completion from retained result data.

View Source
var ErrTransition = errors.New("job state transition rejected")

Functions

This section is empty.

Types

type Access

type Access struct {
	Action Action
	Job    *xasync.Job
	Input  any
}

Access carries a detached public record and, for submit/replay, fresh canonical input. Authorization must revalidate current access and refresh any captured authorization claims on Input. Persisted principal fields alone are not proof.

type Action

type Action string
const (
	Submit  Action = "submit"
	Replay  Action = "replay"
	Inspect Action = "inspect"
)

type Authorizer

type Authorizer func(context.Context, Access) error

type Capture

type Capture struct {
	Plan     *dexec.ReadPlan
	Input    any
	Controls Controls
}

type Config

type Config struct {
	Store     *SQLStore
	Publisher Publisher
	Authorize Authorizer
	TTL       time.Duration
	ErrorTTL  time.Duration
	// Notify receives terminal state only after durable completion write-back.
	// Failure is returned without converting completed work back to ERROR.
	Notify func(context.Context, *xasync.Job) error
}

type Controls

type Controls struct {
	MatchKey string
	JobID    string
	Sync     bool
	Result   bool
}

type Dispatcher

type Dispatcher interface {
	Capture(context.Context, *Record, Submission) (*Capture, error)
	Restore(context.Context, *Record) (Execution, error)
}

Dispatcher adapts jobs to the canonical runtime, never to a second SQL engine.

type Event

type Event struct {
	xasync.Job
	State   string
	Metrics string
	SQL     json.RawMessage
}

Event preserves the original .job document shape, including State and SQL. Public SDK JSON names are case-insensitively readable by original Datly.

type Exchange

type Exchange struct {
	Job    *xasync.Job
	Plan   *dexec.ReadPlan
	Value  any
	Reused bool
}

Exchange is the result of original schedule/sync/inspect orchestration. Value is present only when this invocation actually executed the job.

type Execution

type Execution interface {
	Execute(context.Context) (*ExecutionResult, error)
}

type ExecutionResult

type ExecutionResult struct {
	Value   any
	Outcome *xhandler.Outcome
	Metrics xresponse.Metrics
}

ExecutionResult retains canonical completion evidence without granting transaction handles or completion controls to the job transport.

type InputPolicy

type InputPolicy interface {
	Select(context.Context, any) (Controls, error)
	Owns(*xasync.Job) bool
}

InputPolicy selects scheduling controls from freshly bound canonical input. It is trusted registration metadata, never supplied by a transport client.

type Invocation

type Invocation struct {
	Job  *xasync.Job
	Type xasync.InvocationType
}

Invocation carries the original event-invocation metadata in a fresh replay scope. It is execution metadata, not a global or client-bindable principal.

func CurrentInvocation

func CurrentInvocation(ctx context.Context) (Invocation, bool)

func (Invocation) Context

func (i Invocation) Context(ctx context.Context) context.Context

type Presentation

type Presentation struct {
	Job   *xasync.Job
	Error error
}

Presentation supplies original Datly KindAsync metadata for one durable job. It holds public detached metadata, never raw SourceState or a principal service.

func (*Presentation) DefaultCacheable

func (p *Presentation) DefaultCacheable() bool

func (*Presentation) Kind

func (p *Presentation) Kind() string

func (*Presentation) Locate

func (*Presentation) Priority

func (p *Presentation) Priority() int

func (*Presentation) Value

func (p *Presentation) Value(_ context.Context, _ reflect.Type, name string) (any, bool, error)

type PublicationError

type PublicationError struct {
	ID, EventURL string
	Err          error
}

func (*PublicationError) Error

func (e *PublicationError) Error() string

func (*PublicationError) Unwrap

func (e *PublicationError) Unwrap() error

type Publisher

type Publisher interface {
	Prepare(*Record) error
	Publish(context.Context, *Record) error
}

Publisher owns dispatch-event publication, separate from terminal notification.

type Record

type Record struct {
	xasync.Job
	State    string
	Metrics  string
	SQLQuery string
}

Record adds the original persistence-only columns to the existing public SDK model. State is original session-cache JSON keyed by canonical parameter name. SQLQuery uses original Query/Args JSON objects. There is no StateCodec column.

func (*Record) Event

func (r *Record) Event() (*Event, error)

func (*Record) Public

func (r *Record) Public() (*xasync.Job, error)

func (*Record) QueryScope

func (r *Record) QueryScope() (*sqlxread.QueryScope, error)

QueryScope projects original durable execution metadata to a native SQLX identity guard. SQL is only compared with newly built SQL, never executed here.

type ResultReader

type ResultReader interface {
	ReadResult(context.Context, ResultRequest) (any, error)
}

ResultReader is the canonical dispatcher's optional reader capability. It builds a new authorized typed read, never repeats a mutation or decodes output.

type ResultRequest

type ResultRequest struct {
	Job         *Record
	Source      dexec.ProviderScope
	SourceState string
	Refresh     bool
}

type SQLConfig

type SQLConfig struct {
	DB                   *sql.DB
	Dialect              *info.Dialect
	Table                string
	Dataset              string
	DisableTableCreation bool
}

SQLConfig selects the configured job connector and original-shaped job table. Missing tables are created by SQLX unless DisableTableCreation is set. Existing tables are never altered or rewritten.

type SQLStore

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

SQLStore uses the original native SQLX insert/update and typed read owners. Its database is the configured job connector, never the application Data scope.

func NewSQLStore

func NewSQLStore(ctx context.Context, config SQLConfig) (*SQLStore, error)

func (*SQLStore) Create

func (s *SQLStore) Create(ctx context.Context, r *Record) error

func (*SQLStore) Expire

func (s *SQLStore) Expire(ctx context.Context, now time.Time) (int64, error)

func (*SQLStore) Get

func (s *SQLStore) Get(ctx context.Context, id string) (*Record, error)

func (*SQLStore) Pending

func (s *SQLStore) Pending(ctx context.Context, limit int) ([]*Record, error)

type Scheduled

type Scheduled struct {
	Job  *xasync.Job
	Plan *dexec.ReadPlan
}

type Service

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

func NewService

func NewService(config Config, dispatcher Dispatcher) (*Service, error)

func (*Service) Exchange

func (s *Service) Exchange(ctx context.Context, request Submission) (result *Exchange, failure error)

Exchange prepares external input once, then schedules, runs synchronously or inspects durable state. All execution still goes through Run's canonical claim, replay and completion path. No result replay or second transaction owner exists.

func (*Service) HandleJob

func (s *Service) HandleJob(ctx context.Context, event *Event) (any, error)

HandleJob is the common entry point for local and cloud storage dispatch. Event contents identify a durable record; replay state and route authority are always loaded from DATLY_JOBS, never trusted from the uploaded document.

func (*Service) Republish

func (s *Service) Republish(ctx context.Context, id string) error

Republish retries transport delivery for an authorized pending durable job. It does not reset a running/terminal job or repeat application execution.

func (*Service) Run

func (s *Service) Run(ctx context.Context, id string) (result any, err error)

func (*Service) Schedule

func (s *Service) Schedule(ctx context.Context, request Submission) (*Scheduled, error)

func (*Service) Status

func (s *Service) Status(ctx context.Context, id string) (*xasync.Job, error)

type StateCodec

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

StateCodec bounds original name-keyed state; all field, presence, source-type and replay binding mechanics belong to the native canonical Bindly plan.

func NewStateCodec

func NewStateCodec(contract *registry.RouteInputContract) (*StateCodec, error)

func (*StateCodec) Capture

func (c *StateCodec) Capture(input any) (*bindly.Replay, error)

func (*StateCodec) CaptureSources

func (c *StateCodec) CaptureSources(ctx context.Context, providers []locator.Provider) (*bindly.Replay, error)

CaptureSources delegates raw source authority to the registered native plan.

func (*StateCodec) Decode

func (c *StateCodec) Decode(state string) (*bindly.Replay, error)

func (*StateCodec) Encode

func (c *StateCodec) Encode(input any) (string, error)

func (*StateCodec) State

func (c *StateCodec) State(replay *bindly.Replay) (string, error)

type StatusPresentation

type StatusPresentation struct{ Error error }

StatusPresentation supplies the existing output/status shape for metadata-only replies. Actual handler result finalization remains inside the execution engine.

func (*StatusPresentation) DefaultCacheable

func (p *StatusPresentation) DefaultCacheable() bool

func (*StatusPresentation) Kind

func (p *StatusPresentation) Kind() string

func (*StatusPresentation) Locate

func (*StatusPresentation) Priority

func (p *StatusPresentation) Priority() int

func (*StatusPresentation) Value

func (p *StatusPresentation) Value(_ context.Context, target reflect.Type, name string) (any, bool, error)

type Submission

type Submission struct {
	Job   xasync.Job
	Input any
	// SourceState supplies original parameter-name-keyed raw values for codec
	// inputs or bodies whose authored presence cannot be recovered from a value.
	SourceState string
	Source      dexec.ProviderScope
	Policy      InputPolicy
}

Jump to

Keyboard shortcuts

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