Documentation
¶
Overview ¶
Package jobs owns durable scheduling over the original Datly DATLY_JOBS table.
Index ¶
- Constants
- Variables
- type Access
- type Action
- type Authorizer
- type Capture
- type Config
- type Controls
- type Dispatcher
- type Event
- type Exchange
- type Execution
- type ExecutionResult
- type InputPolicy
- type Invocation
- type Presentation
- type PublicationError
- type Publisher
- type Record
- type ResultReader
- type ResultRequest
- type SQLConfig
- type SQLStore
- type Scheduled
- type Service
- func (s *Service) Exchange(ctx context.Context, request Submission) (result *Exchange, failure error)
- func (s *Service) HandleJob(ctx context.Context, event *Event) (any, error)
- func (s *Service) Republish(ctx context.Context, id string) error
- func (s *Service) Run(ctx context.Context, id string) (result any, err error)
- func (s *Service) Schedule(ctx context.Context, request Submission) (*Scheduled, error)
- func (s *Service) Status(ctx context.Context, id string) (*xasync.Job, error)
- type StateCodec
- func (c *StateCodec) Capture(input any) (*bindly.Replay, error)
- func (c *StateCodec) CaptureSources(ctx context.Context, providers []locator.Provider) (*bindly.Replay, error)
- func (c *StateCodec) Decode(state string) (*bindly.Replay, error)
- func (c *StateCodec) Encode(input any) (string, error)
- func (c *StateCodec) State(replay *bindly.Replay) (string, error)
- type StatusPresentation
- func (p *StatusPresentation) DefaultCacheable() bool
- func (p *StatusPresentation) Kind() string
- func (p *StatusPresentation) Locate(*structology.State) locator.Locator
- func (p *StatusPresentation) Priority() int
- func (p *StatusPresentation) Value(_ context.Context, target reflect.Type, name string) (any, bool, error)
- type Submission
Constants ¶
const MaxInlineState = 63 * 1024
Variables ¶
var ErrCompletionPending = errors.New("job completion requires managed transaction reconciliation")
var ErrInProgress = errors.New("job execution is already in progress")
var ErrJobFailed = errors.New("job execution failed")
var ErrNotFound = errors.New("job not found")
ErrResultUnavailable distinguishes durable completion from retained result data.
var ErrTransition = errors.New("job state transition rejected")
Functions ¶
This section is empty.
Types ¶
type Access ¶
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 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 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 ¶
Event preserves the original .job document shape, including State and SQL. Public SDK JSON names are case-insensitively readable by original Datly.
type Exchange ¶
Exchange is the result of original schedule/sync/inspect orchestration. Value is present only when this invocation actually executed the job.
type ExecutionResult ¶
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)
type Presentation ¶
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 (p *Presentation) Locate(*structology.State) locator.Locator
func (*Presentation) Priority ¶
func (p *Presentation) Priority() int
type PublicationError ¶
func (*PublicationError) Error ¶
func (e *PublicationError) Error() string
func (*PublicationError) Unwrap ¶
func (e *PublicationError) Unwrap() error
type Record ¶
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) 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.
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 ¶
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 ¶
Republish retries transport delivery for an authorized pending durable job. It does not reset a running/terminal job or repeat application execution.
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) CaptureSources ¶
func (c *StateCodec) CaptureSources(ctx context.Context, providers []locator.Provider) (*bindly.Replay, error)
CaptureSources delegates raw source authority to the registered native plan.
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 (p *StatusPresentation) Locate(*structology.State) locator.Locator
func (*StatusPresentation) Priority ¶
func (p *StatusPresentation) Priority() int
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
}