diff --git a/.changeset/20281-job-pull-organization.md b/.changeset/20281-job-pull-organization.md new file mode 100644 index 00000000000..acca3fdf16b --- /dev/null +++ b/.changeset/20281-job-pull-organization.md @@ -0,0 +1,19 @@ +--- +'@objectstack/spec': minor +'@objectstack/runtime': minor +'@objectstack/service-automation': minor +--- + +A job pulls a mapping's connector source by declaration — `pull: { mapping }` — and every job runs as the `organization` it declares (#20281). + +Clause-②: yes (widening) + +- **`JobSchema.pull`** (`@objectstack/spec/system`). A third run form beside `body` and `handler`: `{ mapping: '' }`. On each run the platform pulls that mapping's `connectorSource` and writes the rows through the import runner. It carries no code. The key is refused beside `body` or `handler`, because one of the two run forms would never run. `body` with `handler` stays legal, and the body still wins. A job must now declare one of `body`, `handler` or `pull`. `pull` is closed: an unknown key inside it is refused. +- **`JobSchema.organization`**. The organization a job runs as. It applies to the body's `ctx.api`, to the handler's new `executionContext`, and to the pull's reads and writes. The value shape is the scheduled flow's: a non-empty `sys_organization.id`. A near-miss spelling (`organizationId`, `orgId`, `tenantId`, …) is refused at parse and pointed at the key. +- **`defineStack`, and so `os validate`**, refuses a job whose `pull` names a mapping the stack does not declare, or a mapping with no `connectorSource`. The refusal is the existing `STACK_CROSS_REFERENCE_INVALID` envelope. +- **`IAutomationService.pullConnectorSource`** (`@objectstack/spec/contracts`, with `ConnectorSourcePullRequest`, `ConnectorSourcePullResult` and `ConnectorSourcePullSummary`). The connector sync executor is now on the `automation` service. `@objectstack/service-automation`'s engine serves it from the executor `AutomationServicePlugin` attaches at init (`AutomationEngine.setConnectorPullSource`). A bare engine refuses with `SERVICE_UNAVAILABLE` (503). +- **The job binder** (`@objectstack/runtime`, `scheduleAppArtifactJobs`) schedules a `pull` job on every door: the boot, and `os package install` on install and rehydrate. Each run calls `pullConnectorSource` through the service registry. A refused pull fails the run, and `retryPolicy` applies. A pull whose rows the import runner refused records the run `degraded`, with the counts. A pull naming a mapping the artifact does not carry is not scheduled, and neither is one whose mapping has no `connectorSource`, nor one on a kernel whose `automation` service cannot pull. Each case is logged at `warn` with the reason. `collectJobsWithoutBody` no longer names a `pull` job, so `os package install` does not refuse one. The result gains `pulls` and `missingOrganization`. +- **The organization, judged at bind** by the posture rule scheduled flows use (`resolveScheduledWorkPolicy`). Every run carries `{ isSystem: true, tenantId: }`, or `{ isSystem: true }` for a job that declares none. Under `single` the key is not required. Under `group` it is optional; an undeclared job is scheduled and named once at `warn`, because a tenant-scoped row it writes is refused. Under `isolated`, with package-authored scheduled work switched on, it is **required**. **Action on such a deployment:** declare `organization` on each packaged job, or the job is not scheduled; the error log names the job. Until now such a job was scheduled, and every tenant-scoped write it made was refused at the write. An unrecognized `OS_TENANCY_POSTURE` withholds every job (`scheduled-work-policy-unreadable`) instead of guessing whether a declaration is required. +- **Texts this makes true.** The `mapping.connectorSource` description, the `connector.syncConfig` tombstone prescription and the `connector-sync-keys-retired` upgrade entry said "nothing schedules a pull yet". They now name the `job` `pull` that drives it. + +Nothing that parsed before is refused now. Every new refusal falls on a key that did not exist before this change. diff --git a/content/docs/automation/jobs.mdx b/content/docs/automation/jobs.mdx index 1107a4a6180..c31edc5b3ae 100644 --- a/content/docs/automation/jobs.mdx +++ b/content/docs/automation/jobs.mdx @@ -1,14 +1,16 @@ --- title: Scheduled jobs — cron automation metadata navTitle: Scheduled Jobs -description: Run a TypeScript function on a cron, interval, or one-off schedule — and decide when a job is the right tool instead of a schedule-triggered flow. +description: Run a sandboxed body, a connector pull, or a TypeScript function on a cron, interval, or one-off schedule — and decide when a job is the right tool instead of a schedule-triggered flow. --- A **job** runs work on a schedule. You declare the schedule as metadata; the platform's job service owns the timing, the retries, the per-attempt time limit, -and the run history. The work is either a sandboxed [`body`](#the-job-body) -that travels with the metadata — the preferred form — or, deprecated, the name of -a function in your bundle (`handler`). +and the run history. The work is a sandboxed [`body`](#the-job-body) that travels +with the metadata — the preferred form for code — a [`pull`](#pulling-a-mapping) of +a mapping's connector source, which is no code at all, or, deprecated, the name of a +function in your bundle (`handler`). A job runs as the +[organization it declares](#the-organization-a-job-runs-as). {/* os:check */} ```typescript @@ -39,9 +41,9 @@ writes carry*: | | `job` | `schedule`-type flow | |:---|:---|:---| -| What runs | a sandboxed `body`, or (deprecated) one TypeScript function from `defineStack({ functions })` | a node graph — record operations, `notify`, `http`, approvals, subflows | +| What runs | a sandboxed `body`, a mapping `pull`, or (deprecated) one TypeScript function from `defineStack({ functions })` | a node graph — record operations, `notify`, `http`, approvals, subflows | | Changeable after deploy | **No.** `job` is `allowRuntimeCreate: false` and `allowOrgOverride: false` — there is no "create job" in Studio and no per-tenant fork | Yes — a new flow can be authored through Studio / `PUT /meta` (`allowRuntimeCreate: true`) | -| Identity of its data writes | whatever the handler does with the engine it is given | declared by [`runAs`](/docs/automation/flows) — and a `user` run that resolves no trigger user has its data operations **refused**, so a scheduled flow normally declares `runAs: 'system'` | +| Identity of its data writes | system, in the job's declared [`organization`](#the-organization-a-job-runs-as) | declared by [`runAs`](/docs/automation/flows) — and a `user` run that resolves no trigger user has its data operations **refused**, so a scheduled flow normally declares `runAs: 'system'` | | Retry / time limit | `retryPolicy` + `timeoutMs` on the job, honoured by the job adapter | the flow's own error handling | | Run history | `sys_job` + `sys_job_run` | `sys_automation_run` | @@ -141,7 +143,8 @@ A job's `body` is the same sandboxed JavaScript body hooks and script actions carry — `{ language: 'js', source, capabilities }` — so the work travels with the metadata instead of living in a runtime module only some boots import. When a job declares both, `body` wins; `handler` is deprecated beside it. A job must declare -at least one of the two. +one of `body`, `handler` or [`pull`](#pulling-a-mapping), and `pull` is refused +beside either of the other two. {/* os:check */} ```typescript @@ -173,6 +176,7 @@ export const CloseStaleTasksJob = defineJob({ package whose enabled job has no `body`, or a `body` that does not bind (an expression body, or one carrying `body.timeoutMs`), with `422 VALIDATION_ERROR` and the remedy: give the job a valid `body`, or boot it with `os start --artifact`. + A [`pull`](#pulling-a-mapping) job is data too, and is not refused. Uninstalling a package stops its scheduled jobs at once, and a reinstall whose new version drops a job stops that job. @@ -218,6 +222,82 @@ sweep needs. A body also runs under the sandbox's per-run memory cap (`body.memoryMb`, at most 256), which is one more reason to page rather than load everything at once. +## Pulling a mapping + +A job whose work is "copy the records an external system holds into a local +object" needs no code. Declare the sync on its target — a +[`mapping`](/docs/references/data/mapping) with a `connectorSource` naming the +`rest` or `openapi` connector it reads from — and give the job a `pull` naming +that mapping. The mapping says where the rows come from, how fields map, and how a +pulled row matches a stored one; the job says when. + +{/* os:check */} +```typescript +import { defineJob } from '@objectstack/spec'; + +export const OrdersPullJob = defineJob({ + name: 'orders_pull_hourly', + schedule: { type: 'cron', expression: '0 * * * *', timezone: 'UTC' }, + pull: { mapping: 'orders_pull' }, + retryPolicy: { maxRetries: 2, backoffMs: 60000 }, +}); +``` + +- **A run form of its own.** `pull` is refused beside `body` or `handler`: the + platform binds the pull itself, so code beside it would never run. A body + cannot call a pull; a pull on a schedule is this declaration. +- **Checked when you build.** `defineStack` — and so `os validate` — refuses a + `pull` that names a mapping the stack does not declare, or one with no + `connectorSource`. The binder checks the same reference against the artifact + before it schedules the job, on every door. +- **Each run is one pull.** It calls the connector's read action once, reads one + response, and writes the records through the import runner with the mapping's + `mode` and `upsertKey` — see + [Data sync is defined on the target](https://github.com/objectstack-ai/objectstack/blob/main/packages/spec/docs/SYNC_ARCHITECTURE.md#data-sync-is-defined-on-the-target) + for the watermark and the one-response limit. +- **How a run is recorded.** A pull the platform refuses — the mapping is gone, + the connector is degraded, the upstream answered `ok: false` — rejects, so the + run is `failed` and the `retryPolicy` applies. A pull whose rows the import + runner refused completes as `degraded`, with the counts as its reason; retrying + would refuse the same rows. Otherwise the run is `success`, including a pull + that found nothing new. + +## The organization a job runs as + +A scheduled run has no session to inherit an organization from. A job declares the +one it runs in: + +{/* os:check */} +```typescript +import { defineJob } from '@objectstack/spec'; + +export const PlantSweepJob = defineJob({ + name: 'plant_a_nightly_sweep', + schedule: { type: 'cron', expression: '0 2 * * *', timezone: 'UTC' }, + organization: 'org_plant_a', + body: { + language: 'js', + source: "await ctx.api.object('task').find({ where: { status: 'open' }, limit: 100 });", + capabilities: ['api.read'], + }, +}); +``` + +Every run form runs as it: a `body`'s `ctx.api`, a `pull`'s reads and writes, and +the `executionContext` a `handler` is handed are all system access carrying that +organization, so a tenant-scoped row the job writes is stamped with it. Whether the +key is required is the same deployment-posture rule +[time-triggered flows](/docs/automation/flows) use, read when the job is scheduled: + +| Tenancy posture (with packaged scheduled work switched on) | `organization` | A job that declares none | +|:---|:---|:---| +| `single` | not required | runs with none; the install's one organization is resolved beneath each write | +| `group` | optional | is scheduled, and a tenant-scoped row it writes is **refused** — the boot names such jobs once, at `warn` | +| `isolated` | **required** | is **not scheduled**, logged at `error` with the remedy | + +No organization is ever chosen for a job that declares none. Work wanted in +several organizations is one job per organization. + ## The handler, and the ways it can fail to be one `handler` must match a key of `defineStack({ functions })`. At `kernel:ready` @@ -248,6 +328,7 @@ At run time the handler is invoked with a `JobHandlerContext` | `bundle` | the application's metadata bundle — declarations, not a data handle | | `ql` | the live ObjectQL engine — the same handle `defineStack({ onEnable })` receives | | `logger` | the platform logger, so a job's diagnostics are not `console` output | +| `executionContext` | the context the job runs as — `{ isSystem: true, tenantId }` for a job that declares an [`organization`](#the-organization-a-job-runs-as), else `{ isSystem: true }`. `ql` is the raw engine, so pass it as each call's `context` to write as that organization | ```typescript import type { JobHandlerContext } from '@objectstack/runtime'; diff --git a/content/docs/permissions/system-context.mdx b/content/docs/permissions/system-context.mdx index 6c1ad8a6859..b8083cac832 100644 --- a/content/docs/permissions/system-context.mdx +++ b/content/docs/permissions/system-context.mdx @@ -353,7 +353,7 @@ still holds equal to the census on every pull request: | — in tests | 1013 | — | | — in non-test sources | 798 | — | | Appearances of the bare identifier `isSystem` in non-test sources | 813 | — | -| — parsed as a declaration | 23 | ✅ | +| — parsed as a declaration | 25 | ✅ | | — parsed as an object-literal / type key (producers and option objects) | 310 | — | | — parsed as a property **read** | 120 | ✅ | | — parsed in some other syntactic position (a local, a cast, a conditional) | 9 | ✅ | diff --git a/content/docs/references/data/mapping.mdx b/content/docs/references/data/mapping.mdx index 09328256be3..bad43bad443 100644 --- a/content/docs/references/data/mapping.mdx +++ b/content/docs/references/data/mapping.mdx @@ -49,7 +49,7 @@ const result = ImportFieldMappingSchema.parse(data); | **fieldMapping** | `{ source: string \| string[]; target: string \| string[]; transform?: Enum<'none' \| 'constant' \| 'lookup' \| 'split' \| 'join' \| 'javascript' \| 'map'>; params?: object }[]` | ✅ | | | **mode** | `Enum<'insert' \| 'update' \| 'upsert'>` | optional (default: `"insert"`) | | | **upsertKey** | `string[]` | optional | Fields to match for upsert (e.g. email) | -| **connectorSource** | `{ connector: string; action: string; input?: Record; recordsPath?: string; … }` | optional | Pull binding: the rest/openapi connector this mapping pulls rows from (one-way, full or timestamp-incremental; a `job` sets the cadence). Pulled when a job drives it; nothing schedules it yet, so the binding alone moves no rows — schedule the pull with a `job` once a job can drive one | +| **connectorSource** | `{ connector: string; action: string; input?: Record; recordsPath?: string; … }` | optional | Pull binding: the rest/openapi connector this mapping pulls rows from (one-way, full or timestamp-incremental; a `job` sets the cadence). Pulled when a job drives it — a `job` whose `pull: { mapping }` names this mapping, on the job's schedule; the binding alone moves no rows | | **_lock** | `Enum<'none' \| 'no-overlay' \| 'no-delete' \| 'full'>` | optional | Item-level lock — controls overlay & delete (ADR-0010). | | **_lockReason** | `string` | optional | Human-readable reason shown when a write is refused by _lock. | | **_lockSource** | `Enum<'artifact' \| 'package' \| 'env-forced'>` | optional | Layer that set _lock (artifact \| package \| env-forced). | diff --git a/content/docs/references/integration/connector.mdx b/content/docs/references/integration/connector.mdx index 02b3d7a0962..c8f93c87e0f 100644 --- a/content/docs/references/integration/connector.mdx +++ b/content/docs/references/integration/connector.mdx @@ -191,7 +191,7 @@ const result = ConnectorSchema.parse(data); | **auth** | `{ type: 'none' } \| { type: 'bearer'; credentialRef: string } \| { type: 'api-key'; credentialRef: string; headerName?: string; paramName?: string } \| { type: 'basic'; username: string; credentialRef: string }` | optional | Declarative instance auth — references credentials via `credentialRef` (resolved at boot), never inline secrets. Requires `provider` (ADR-0097). | | **actions** | `{ key: string; label: string; description?: string; inputSchema?: Record; … }[]` | optional | | | **triggers** | `never` | optional | [REMOVED] `connector.triggers` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — a connector trigger never started anything: `AutomationEngine.registerConnector` registers a connector's actions only, no polling loop read `intervalSeconds` (or the `interval` spelling it was renamed from), and no receiver was driven by a `webhook` trigger. Delete the key; the `ConnectorTrigger` shape leaves with it. To start work from an external system, write a flow that calls the connector's action in a `connector_action` node: for an external event, an `api` flow that the event's sender calls; for a scheduled pull, a `schedule` flow. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | -| **syncConfig** | `never` | optional | [REMOVED] `connector.syncConfig` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — no engine ever ran a connector-attached sync: nothing read `strategy`, `direction`, `realtimeSync`, `timestampField`, `conflictResolution`, `batchSize`, `deleteMode` or `filters`, so the `latest_wins` and `soft_delete` defaults resolved and deleted nothing. Delete the key; the `DataSyncConfig` shape leaves with it. A sync is defined on its TARGET: a `mapping` (`targetObject`, `fieldMapping`, `mode`, `upsertKey`) whose `connectorSource` names the `rest` or `openapi` connector it pulls from, the read action and an optional timestamp `watermark`, with a `job` for the cadence. Its pull runs when a `job` drives it; nothing schedules one yet, so the binding alone moves no rows. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | +| **syncConfig** | `never` | optional | [REMOVED] `connector.syncConfig` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — no engine ever ran a connector-attached sync: nothing read `strategy`, `direction`, `realtimeSync`, `timestampField`, `conflictResolution`, `batchSize`, `deleteMode` or `filters`, so the `latest_wins` and `soft_delete` defaults resolved and deleted nothing. Delete the key; the `DataSyncConfig` shape leaves with it. A sync is defined on its TARGET: a `mapping` (`targetObject`, `fieldMapping`, `mode`, `upsertKey`) whose `connectorSource` names the `rest` or `openapi` connector it pulls from, the read action and an optional timestamp `watermark`, with a `job` for the cadence. Its pull runs when a `job` drives it — a `job` whose `pull: { mapping }` names that mapping; the binding alone moves no rows. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | | **fieldMappings** | `never` | optional | [REMOVED] `connector.fieldMappings` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — no engine ever moved a value through a connector field mapping: nothing read `source`, `target`, `defaultValue`, `dataType`, `required` or `syncMode`. Delete the key; the `ConnectorFieldMapping` shape leaves with it. Map fields on the sync's TARGET instead: a `mapping`'s `fieldMapping` (`source` → `target`, with a `transform` the import path executes), which its `connectorSource` pulls through. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | | **webhooks** | `never` | optional | [REMOVED] `connector.webhooks` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — a webhook nested inside a connector was never registered as a `webhook` item, so it was never materialized into `sys_webhook` and never delivered, and nothing emits the connector events its `events` list could name (`sync.completed`, `auth.expired` and the rest). Delete the key; the nested shape leaves with it (`WebhookConfig`, `WebhookEvent`, `WebhookSignatureAlgorithm`). To have a webhook actually sent, declare it in the stack's top-level `webhooks:` collection, which is materialized into `sys_webhook` and delivered on record events — note that doing so STARTS deliveries this connector never made. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | | **rateLimitConfig** | `never` | optional | [REMOVED] `connector.rateLimitConfig` was removed in @objectstack/spec 17.0.0 (ADR-0049 D2) — the entire shape is gone, not just this key: `ConnectorRateLimitConfig` and its `RateLimitStrategy` enum were removed with it, because no outbound rate-limiting engine ever existed. The platform's only token bucket (runtime `security/rate-limit.ts`) throttles INBOUND requests to us; nothing throttled the calls a connector makes out, so every knob here was inert while reading like a configured cap. Delete the key. Do NOT substitute `shared` `RateLimitConfig` — that is the inbound limiter and would cap the wrong direction; until an outbound throttle exists, rate-limit at the connector provider or upstream gateway. Run `os migrate meta --from 16` to list the mechanical edits for existing sources; apply them by hand. | @@ -488,7 +488,7 @@ Connector type | **auth** | `{ type: 'none' } \| { type: 'bearer'; credentialRef: string } \| { type: 'api-key'; credentialRef: string; headerName?: string; paramName?: string } \| { type: 'basic'; username: string; credentialRef: string }` | optional | Declarative instance auth — references credentials via `credentialRef` (resolved at boot), never inline secrets. Requires `provider` (ADR-0097). | | **actions** | `{ key: string; label: string; description?: string; inputSchema?: Record; … }[]` | optional | | | **triggers** | `never` | optional | [REMOVED] `connector.triggers` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — a connector trigger never started anything: `AutomationEngine.registerConnector` registers a connector's actions only, no polling loop read `intervalSeconds` (or the `interval` spelling it was renamed from), and no receiver was driven by a `webhook` trigger. Delete the key; the `ConnectorTrigger` shape leaves with it. To start work from an external system, write a flow that calls the connector's action in a `connector_action` node: for an external event, an `api` flow that the event's sender calls; for a scheduled pull, a `schedule` flow. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | -| **syncConfig** | `never` | optional | [REMOVED] `connector.syncConfig` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — no engine ever ran a connector-attached sync: nothing read `strategy`, `direction`, `realtimeSync`, `timestampField`, `conflictResolution`, `batchSize`, `deleteMode` or `filters`, so the `latest_wins` and `soft_delete` defaults resolved and deleted nothing. Delete the key; the `DataSyncConfig` shape leaves with it. A sync is defined on its TARGET: a `mapping` (`targetObject`, `fieldMapping`, `mode`, `upsertKey`) whose `connectorSource` names the `rest` or `openapi` connector it pulls from, the read action and an optional timestamp `watermark`, with a `job` for the cadence. Its pull runs when a `job` drives it; nothing schedules one yet, so the binding alone moves no rows. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | +| **syncConfig** | `never` | optional | [REMOVED] `connector.syncConfig` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — no engine ever ran a connector-attached sync: nothing read `strategy`, `direction`, `realtimeSync`, `timestampField`, `conflictResolution`, `batchSize`, `deleteMode` or `filters`, so the `latest_wins` and `soft_delete` defaults resolved and deleted nothing. Delete the key; the `DataSyncConfig` shape leaves with it. A sync is defined on its TARGET: a `mapping` (`targetObject`, `fieldMapping`, `mode`, `upsertKey`) whose `connectorSource` names the `rest` or `openapi` connector it pulls from, the read action and an optional timestamp `watermark`, with a `job` for the cadence. Its pull runs when a `job` drives it — a `job` whose `pull: { mapping }` names that mapping; the binding alone moves no rows. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | | **fieldMappings** | `never` | optional | [REMOVED] `connector.fieldMappings` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — no engine ever moved a value through a connector field mapping: nothing read `source`, `target`, `defaultValue`, `dataType`, `required` or `syncMode`. Delete the key; the `ConnectorFieldMapping` shape leaves with it. Map fields on the sync's TARGET instead: a `mapping`'s `fieldMapping` (`source` → `target`, with a `transform` the import path executes), which its `connectorSource` pulls through. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | | **webhooks** | `never` | optional | [REMOVED] `connector.webhooks` was removed in @objectstack/spec 17 (ADR-0049 enforce-or-remove) — a webhook nested inside a connector was never registered as a `webhook` item, so it was never materialized into `sys_webhook` and never delivered, and nothing emits the connector events its `events` list could name (`sync.completed`, `auth.expired` and the rest). Delete the key; the nested shape leaves with it (`WebhookConfig`, `WebhookEvent`, `WebhookSignatureAlgorithm`). To have a webhook actually sent, declare it in the stack's top-level `webhooks:` collection, which is materialized into `sys_webhook` and delivered on record events — note that doing so STARTS deliveries this connector never made. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | | **rateLimitConfig** | `never` | optional | [REMOVED] `connector.rateLimitConfig` was removed in @objectstack/spec 17.0.0 (ADR-0049 D2) — the entire shape is gone, not just this key: `ConnectorRateLimitConfig` and its `RateLimitStrategy` enum were removed with it, because no outbound rate-limiting engine ever existed. The platform's only token bucket (runtime `security/rate-limit.ts`) throttles INBOUND requests to us; nothing throttled the calls a connector makes out, so every knob here was inert while reading like a configured cap. Delete the key. Do NOT substitute `shared` `RateLimitConfig` — that is the inbound limiter and would cap the wrong direction; until an outbound throttle exists, rate-limit at the connector provider or upstream gateway. Run `os migrate meta --from 16` to list the mechanical edits for existing sources; apply them by hand. | diff --git a/content/docs/references/system/job.mdx b/content/docs/references/system/job.mdx index cf802d57bf7..bfbd69a11e2 100644 --- a/content/docs/references/system/job.mdx +++ b/content/docs/references/system/job.mdx @@ -57,8 +57,10 @@ const result = CronScheduleSchema.parse(data); | **label** | `string` | optional | Human-readable label | | **description** | `string` | optional | Job description / purpose | | **schedule** | `{ type: 'cron'; expression: string \| object; timezone?: string } \| { type: 'interval'; intervalMs: integer } \| { type: 'once'; at: string }` | ✅ | Job schedule configuration | -| **handler** | `string` | optional | Handler function name (must match a key in `defineStack({ functions })`) — DEPRECATED, prefer `body`. When both are present `body` wins; a job must declare one of the two. | -| **body** | `{ language: 'js'; source: string; capabilities?: Enum<'api.read' \| 'api.write' \| 'api.transaction' \| 'crypto.uuid' \| 'log'>[]; timeoutMs?: integer; … }` | optional | Job body — a sandboxed JS (L2) body, the same shape hooks and actions use; an expression (L1) body is refused, because a job runs for its effects and an expression has none. Preferred over `handler`: when both are present `body` wins. It runs in the QuickJS sandbox with no module scope (no imports, no helpers or constants from the surrounding file): it reaches data only through `ctx.api` under its declared `capabilities` (`api.read` / `api.write` / `api.transaction`) and logs through `ctx.log` (`log`); the in-process handler context (`ql`, `logger`, `bundle`) does not exist there. Its time limit is the job's `timeoutMs` (see there): long-running work declares a `timeoutMs` that covers it, or splits into bounded runs that each finish within it. Every door that brings an artifact in schedules a job's `body` — the boot, and `os package install` on install and on every restart — while a `handler` is code that travels only in the artifact's runtime module and runs only on a boot that loads it (a config, or `os start --artifact`); `os package install` therefore refuses an enabled job with no `body`. | +| **handler** | `string` | optional | Handler function name (must match a key in `defineStack({ functions })`) — DEPRECATED, prefer `body`. When both are present `body` wins; refused beside `pull`. A job must declare one of `body`, `handler` or `pull`. | +| **body** | `{ language: 'js'; source: string; capabilities?: Enum<'api.read' \| 'api.write' \| 'api.transaction' \| 'crypto.uuid' \| 'log'>[]; timeoutMs?: integer; … }` | optional | Job body — a sandboxed JS (L2) body, the same shape hooks and actions use; an expression (L1) body is refused, because a job runs for its effects and an expression has none. Preferred over `handler`: when both are present `body` wins. It runs in the QuickJS sandbox with no module scope (no imports, no helpers or constants from the surrounding file): it reaches data only through `ctx.api` under its declared `capabilities` (`api.read` / `api.write` / `api.transaction`) and logs through `ctx.log` (`log`); the in-process handler context (`ql`, `logger`, `bundle`) does not exist there. Its time limit is the job's `timeoutMs` (see there): long-running work declares a `timeoutMs` that covers it, or splits into bounded runs that each finish within it. Every door that brings an artifact in schedules a job's `body` — the boot, and `os package install` on install and on every restart — while a `handler` is code that travels only in the artifact's runtime module and runs only on a boot that loads it (a config, or `os start --artifact`); `os package install` therefore refuses an enabled job with no `body` (a `pull` job excepted: it is data too). Refused beside `pull`. | +| **pull** | `{ mapping: string }` | optional | Pull run form: on each run the platform pulls the named mapping's `connectorSource` (one action call, one response) and writes the rows through the import runner — no code. A refused pull records the run `failed` (retried per `retryPolicy`); a pull whose rows the import runner refused records it `degraded`. Data like `body`, so every door schedules it. Refused beside `body` or `handler`. | +| **organization** | `string` | optional | Organization id (sys_organization.id) this job runs as — its `body`'s `ctx.api`, the execution context its `handler` is handed, and its `pull`'s reads and writes alike, as a system run carrying that organization. A scheduled run has no session to inherit one from. Judged at bind by the posture rule scheduled flows use: required under the isolated tenancy posture (a job that declares none is not scheduled); optional under group (undeclared, the run carries no organization and a tenant-scoped write it makes is refused); not required under single (the install's one organization is resolved beneath each write). Where declared, it is the organization the run acts as on every posture. | | **retryPolicy** | `{ maxRetries?: integer; backoffMs?: integer; backoffMultiplier?: number; maxRetryDelayMs?: integer; … }` | optional | Retry policy: failed runs (including timeouts) are retried with exponential backoff (delay = min(backoffMs * backoffMultiplier^(retry-1), maxRetryDelayMs), optionally jittered) up to maxRetries retries after the initial attempt. Omit the block for a single attempt; declaring it without `maxRetries` also means no retry since 17.0.0 — state a count to opt in. | | **timeoutMs** | `integer` | optional | Per-attempt time limit in milliseconds; an over-limit run is recorded with execution status "timeout". A `handler` run is abandoned, not forcibly cancelled. For a job with a `body` this is the ONE time limit: one attempt is one sandbox invocation, the runtime bounds that invocation by this value, and the body shape's own `timeoutMs` (capped at 30000 for hooks and actions) is refused on a job — so this key, which has no such cap, is where long-running work states how long it needs. Omit for no per-attempt limit; a `body` run is then still bounded by the sandbox's own default invocation limits. | | **timeout** | `never` | optional | [REMOVED] `job.timeout` was removed in @objectstack/spec 17 — its unit (milliseconds) lived only in the description while the sibling `retryPolicy.backoffMs` spells its own, so the same number read as two conventions on one surface. Rename the key to `timeoutMs`; the value (milliseconds) is unchanged. Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand. | @@ -103,6 +105,12 @@ const result = CronScheduleSchema.parse(data); | **timeoutMs** | `integer` | optional | Per-invocation timeout (ms) | | **memoryMb** | `integer` | optional | Per-invocation memory cap (MB) | +### Nested Shape: `Job.pull` + +| Property | Type | Required | Description | +| :--- | :--- | :--- | :--- | +| **mapping** | `string` | ✅ | Name of the `mapping` whose `connectorSource` this job pulls on its schedule — a mapping the same stack declares, with a `connectorSource` (refused at `defineStack` / `os validate` otherwise). The mapping says where the rows come from, the field map, the write mode and the match key; the job says when. | + ### Nested Shape: `Job.retryPolicy` | Property | Type | Required | Description | diff --git a/docs/audits/2026-07-unknown-key-strictness-ledger.counts/system.md b/docs/audits/2026-07-unknown-key-strictness-ledger.counts/system.md index 84cee18fd58..0d9366a9ebe 100644 --- a/docs/audits/2026-07-unknown-key-strictness-ledger.counts/system.md +++ b/docs/audits/2026-07-unknown-key-strictness-ledger.counts/system.md @@ -19,4 +19,4 @@ hand-patch a number here** — fix the code or the verdict and regenerate. | Dir | Sites | |---|---| -| `system/` | 353 | +| `system/` | 354 | diff --git a/packages/lint/src/authoring-rules.ts b/packages/lint/src/authoring-rules.ts index b3ce5aa402d..49aa36225a9 100644 --- a/packages/lint/src/authoring-rules.ts +++ b/packages/lint/src/authoring-rules.ts @@ -1585,8 +1585,8 @@ export const AUTHORING_RULES: readonly AuthoringRule[] = [ // sync binding) went `live` and, since #21127, carries no `authorWarn` — a // warned `live` row made this rule throw instead of warn, so a `mapping` // write authoring it got an `authoring-rule-threw` advisory and `os - // validate` / `os lint` exited 1. The scheduling caveat (nothing schedules - // a pull until the `job` stage lands) is on the key's description, and + // validate` / `os lint` exited 1. The scheduling caveat (a pull runs only + // when a `job`'s `pull` names the mapping) is on the key's description, and // `check:liveness` refuses a warned `live` row. That is the ruled end // state, not a half-landing: the ruling dispatched the wiring and ⛔ no // ledger population («the empty warn maps stay empty until a real property needs a row — zero pull, the wiring is diff --git a/packages/runtime/src/app-artifact-handlers.job-pull.test.ts b/packages/runtime/src/app-artifact-handlers.job-pull.test.ts new file mode 100644 index 00000000000..2223ac510a1 --- /dev/null +++ b/packages/runtime/src/app-artifact-handlers.job-pull.test.ts @@ -0,0 +1,365 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +/** + * #20281 stage ③ — "the `job` driving it", built to rulings Q1-B + Q2-O1. + * + * Q1-B: a job drives a connector pull BY DECLARATION — `JobSchema.pull: + * { mapping }`, a third run form the ONE binder binds (`scheduleAppArtifactJobs`) + * by calling the `automation` service's contract method, + * `IAutomationService.pullConnectorSource`, resolved through the service + * registry. Pinned here: + * + * - a pull job is scheduled, and each run calls the contract method with the + * mapping and the job's execution context; + * - the outcome is mapped once, by the binder: a refused pull REJECTS (the job + * service's `failed`, the retry trigger), a pull with refused rows resolves + * `degraded` with the counts, anything else `completed`; + * - a pull that does not bind (a mapping the artifact does not declare, one + * with no `connectorSource`, code beside it) and a composition with no pull + * door are NOT scheduled, and the reason is said; + * - `collectJobsWithoutBody` never names a pull job — it is data. + * + * Q2-O1: a job declares the organization it runs as, judged at bind by the + * scheduled flows' posture rule (`resolveScheduledWorkPolicy`). Pinned here: + * + * - every form runs as `{ isSystem: true, tenantId }` — the pull's `context`, + * the body's `ctx.api` envelope (against the REAL QuickJS sandbox), the + * handler's `executionContext` — and as `{ isSystem: true }` with none; + * - `isolated` (switch on): a job declaring none is NOT scheduled, at `error`; + * - `group`: an undeclared job is scheduled and named once at `warn`; + * - `single`: an undeclared job is scheduled, silently; + * - an unreadable posture fails closed; with the switch OFF it is never read. + */ + +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import type { PluginContext } from '@objectstack/core'; +import { SCHEDULED_WORK_ENV } from '@objectstack/types'; +import { collectJobsWithoutBody, scheduleAppArtifactJobs } from './app-artifact-handlers.js'; +import { withScheduledWorkOn } from './scheduled-work.test-support.js'; + +withScheduledWorkOn(); + +/** Run each test of a suite under `OS_TENANCY_POSTURE=`, restoring the previous value. */ +function withPosture(posture: string | undefined): void { + let prior: string | undefined; + beforeEach(() => { + prior = process.env.OS_TENANCY_POSTURE; + if (posture === undefined) delete process.env.OS_TENANCY_POSTURE; + else process.env.OS_TENANCY_POSTURE = posture; + }); + afterEach(() => { + if (prior === undefined) delete process.env.OS_TENANCY_POSTURE; + else process.env.OS_TENANCY_POSTURE = prior; + }); +} + +const APP_ID = 'com.example.syncapp'; +const INTERVAL = { type: 'interval', intervalMs: 60000 }; +const MAPPING = { + name: 'orders_pull', + targetObject: 'order', + fieldMapping: [{ source: 'id', target: 'external_id' }], + mode: 'upsert', + upsertKey: ['external_id'], + connectorSource: { connector: 'orders_api', action: 'request' }, +}; +const PULL_JOB = { name: 'orders_pull_hourly', schedule: INTERVAL, pull: { mapping: 'orders_pull' } }; + +/** A pull result as the service answers it. */ +const result = (summary: Partial> = {}, pulled = 3) => ({ + mapping: 'orders_pull', + targetObject: 'order', + connector: 'orders_api', + action: 'request', + pulled, + summary: { total: pulled, processed: pulled, created: pulled, updated: 0, skipped: 0, errors: 0, ok: pulled, cancelled: false, ...summary }, +}); + +function harness(opts: { automation?: unknown } = {}) { + const writes: Array<{ object: string; data: unknown; context: unknown }> = []; + const ql = { + createContext: (context: unknown) => ({ + object: (object: string) => ({ + insert: async (data: Record) => { + writes.push({ object, data, context }); + return { id: `r${writes.length}`, ...data }; + }, + }), + }), + }; + const scheduled = new Map Promise }>(); + const jobService = { + schedule: async (name: string, _schedule: unknown, run: (c: any) => Promise) => { scheduled.set(name, { run }); }, + cancel: async (name: string) => { scheduled.delete(name); }, + trigger: async () => undefined, + }; + const pulls: unknown[] = []; + const defaultAutomation = { + pullConnectorSource: vi.fn(async (request: unknown) => { pulls.push(request); return result(); }), + }; + const services: Record = { + job: jobService, + automation: 'automation' in opts ? opts.automation : defaultAutomation, + }; + const logger = { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn() }; + const ctx = { + logger, + getService: (name: string) => { + if (services[name] !== undefined) return services[name]; + throw new Error(`no ${name}`); + }, + } as unknown as PluginContext; + const schedule = (jobs: unknown[], extra: Record = {}) => + scheduleAppArtifactJobs(ctx, { id: APP_ID, version: '0.1.0', type: 'app', jobs, mappings: [MAPPING], ...extra }, { + appId: APP_ID, ql: ql as any, source: 'Test', + }); + const run = (name: string) => scheduled.get(name)!.run({ jobId: name }); + const said = (level: 'warn' | 'error' | 'info') => logger[level].mock.calls.map((c) => String(c[0])); + return { writes, scheduled, services, pulls, defaultAutomation, logger, schedule, run, said }; +} + +describe('#20281 stage ③ (Q1-B): a job pulls a mapping by declaration, through the automation service', () => { + withPosture(undefined); + + it('schedules a pull job, and each run calls pullConnectorSource with the mapping and the job\'s context', async () => { + const h = harness(); + + const out = await h.schedule([PULL_JOB]); + expect(out.pulls).toEqual(['orders_pull_hourly']); + expect(out.bodies).toEqual([]); + expect(out.handlers).toEqual([]); + expect(out.notScheduled).toEqual([]); + expect([...h.scheduled.keys()]).toEqual(['orders_pull_hourly']); + + const outcome = await h.run('orders_pull_hourly'); + expect(h.pulls).toEqual([{ mapping: 'orders_pull', context: { isSystem: true } }]); + expect(outcome).toEqual({ outcome: 'completed' }); + }); + + it('a pull whose rows the import runner refused resolves degraded, with the counts as its reason', async () => { + const h = harness(); + h.defaultAutomation.pullConnectorSource.mockResolvedValueOnce(result({ created: 2, ok: 2, errors: 3 }, 5)); + await h.schedule([PULL_JOB]); + + const outcome = (await h.run('orders_pull_hourly')) as { outcome: string; reason?: string }; + expect(outcome.outcome).toBe('degraded'); + expect(outcome.reason).toContain('3 of 5'); + }); + + it('a pull that carried no records completes — nothing new is not degraded', async () => { + const h = harness(); + h.defaultAutomation.pullConnectorSource.mockResolvedValueOnce(result({ total: 0, processed: 0, created: 0, ok: 0 }, 0)); + await h.schedule([PULL_JOB]); + + expect(await h.run('orders_pull_hourly')).toEqual({ outcome: 'completed' }); + }); + + it('a refused pull REJECTS the run — the job service records failed and the retry policy applies', async () => { + const h = harness(); + const refusal = Object.assign(new Error('Mapping "orders_pull": orders_api.request answered ok:false'), { + code: 'EXTERNAL_SERVICE_ERROR', status: 502, reason: 'upstream_not_ok', + }); + h.defaultAutomation.pullConnectorSource.mockRejectedValueOnce(refusal); + await h.schedule([PULL_JOB]); + + await expect(h.run('orders_pull_hourly')).rejects.toMatchObject({ code: 'EXTERNAL_SERVICE_ERROR', status: 502 }); + }); + + it('a pull naming a mapping the artifact does not declare is NOT scheduled, and the warn names `pull.mapping`', async () => { + const h = harness(); + + const out = await h.schedule([{ ...PULL_JOB, pull: { mapping: 'orders_pul' } }]); + expect(out.pulls).toEqual([]); + expect(out.notScheduled).toEqual(['orders_pull_hourly']); + expect(h.scheduled.size).toBe(0); + expect(h.said('warn').some((m) => m.includes('pull.mapping') && m.includes("'orders_pul'"))).toBe(true); + expect(h.defaultAutomation.pullConnectorSource).not.toHaveBeenCalled(); + }); + + it('a pull whose mapping declares no connectorSource is NOT scheduled — there is nothing to pull', async () => { + const h = harness(); + const { connectorSource: _dropped, ...importOnly } = MAPPING; + + const out = await h.schedule([PULL_JOB], { mappings: [importOnly] }); + expect(out.notScheduled).toEqual(['orders_pull_hourly']); + expect(h.said('warn').some((m) => m.includes('connectorSource'))).toBe(true); + }); + + it('a pull beside a body is NOT scheduled — neither run form runs', async () => { + const h = harness(); + const body = { language: 'js', capabilities: ['api.write'], source: "await ctx.api.object('order').insert({});" }; + + const out = await h.schedule([{ ...PULL_JOB, body }]); + expect(out.notScheduled).toEqual(['orders_pull_hourly']); + expect(out.bodies).toEqual([]); + expect(h.scheduled.size).toBe(0); + }); + + it('a pull resolves a mapping a sibling package of the same artifact declares (ADR-0130 D4)', async () => { + const h = harness(); + const out = await scheduleAppArtifactJobs( + { logger: h.logger, getService: (n: string) => h.services[n] } as unknown as PluginContext, + { + manifest: { id: 'com.example.maps', name: 'Maps', version: '0.1.0', type: 'app' }, + packages: [ + { manifest: { id: 'com.example.maps', name: 'Maps', version: '0.1.0', type: 'app', mappings: [MAPPING] } }, + { manifest: { id: 'com.example.jobs', name: 'Jobs', version: '0.1.0', type: 'module', jobs: [PULL_JOB] } }, + ], + }, + { appId: APP_ID, ql: undefined, source: 'Test' }, + ); + expect(out.pulls).toEqual(['orders_pull_hourly']); + }); + + it('with no automation service serving pullConnectorSource the pull job is NOT scheduled, and the composition remedy is said', async () => { + const h = harness({ automation: { execute: async () => ({ success: true }) } }); + + const out = await h.schedule([PULL_JOB]); + expect(out.notScheduled).toEqual(['orders_pull_hourly']); + expect(h.scheduled.size).toBe(0); + expect(h.said('warn').some((m) => m.includes('pullConnectorSource') && m.includes('AutomationServicePlugin'))).toBe(true); + }); + + it('the automation service is resolved on every run, through the registry — not the instance seen at bind', async () => { + const h = harness(); + await h.schedule([PULL_JOB]); + const replacement = { pullConnectorSource: vi.fn(async () => result()) }; + h.services.automation = replacement; + + await h.run('orders_pull_hourly'); + expect(replacement.pullConnectorSource).toHaveBeenCalledTimes(1); + expect(h.defaultAutomation.pullConnectorSource).not.toHaveBeenCalled(); + }); + + it('collectJobsWithoutBody never names a pull job — the pull is data, like a body (a handler job beside it is, control)', () => { + const named = collectJobsWithoutBody({ + id: APP_ID, + jobs: [PULL_JOB, { ...PULL_JOB, name: 'unbound_pull', pull: { mapping: 'nope' } }, { name: 'fn_job', schedule: INTERVAL, handler: 'sweep' }], + mappings: [MAPPING], + }); + expect(named.map((j) => j.name)).toEqual(['fn_job']); + }); +}); + +describe('#20281 stage ③ (Q2-O1): every run form runs as the job\'s declared organization', () => { + withPosture('single'); + + it('a pull job\'s context carries the organization as tenantId', async () => { + const h = harness(); + await h.schedule([{ ...PULL_JOB, organization: 'org_a' }]); + + await h.run('orders_pull_hourly'); + expect(h.pulls).toEqual([{ mapping: 'orders_pull', context: { isSystem: true, tenantId: 'org_a' } }]); + }); + + it('a body job\'s ctx.api write runs under { isSystem: true, tenantId } — the real QuickJS sandbox', async () => { + const h = harness(); + const body = { language: 'js', capabilities: ['api.write'], source: "await ctx.api.object('tick').insert({ name: 'tick' });" }; + await h.schedule([ + { name: 'scoped_body', schedule: INTERVAL, body, organization: 'org_a' }, + { name: 'plain_body', schedule: INTERVAL, body }, + ]); + + await h.run('scoped_body'); + await h.run('plain_body'); + expect(h.writes.map((w) => w.context)).toEqual([{ isSystem: true, tenantId: 'org_a' }, { isSystem: true }]); + }); + + it('a handler job is handed the same envelope as executionContext; ql stays the raw engine', async () => { + const h = harness(); + const seen: unknown[] = []; + const sweep = async (jobCtx: { executionContext?: unknown; ql: unknown }) => { seen.push([jobCtx.executionContext, typeof jobCtx.ql]); }; + await h.schedule( + [{ name: 'scoped_fn', schedule: INTERVAL, handler: 'sweep', organization: 'org_a' }, { name: 'plain_fn', schedule: INTERVAL, handler: 'sweep' }], + { functions: { sweep } }, + ); + + await h.run('scoped_fn'); + await h.run('plain_fn'); + expect(seen).toEqual([[{ isSystem: true, tenantId: 'org_a' }, 'object'], [{ isSystem: true }, 'object']]); + }); + + it('under single an undeclared job is scheduled and nothing is said about its organization', async () => { + const h = harness(); + const out = await h.schedule([PULL_JOB]); + expect(out.pulls).toEqual(['orders_pull_hourly']); + expect(out.missingOrganization).toEqual([]); + expect(h.said('warn').some((m) => m.includes('NO acting organization'))).toBe(false); + }); +}); + +describe('#20281 stage ③ (Q2-O1): under isolated, a job that declares no organization is not scheduled', () => { + withPosture('isolated'); + + it('every run form without an organization is refused at error, naming the job and the key; a declared one is scheduled', async () => { + const h = harness(); + const body = { language: 'js', capabilities: ['api.read'], source: "await ctx.api.object('tick').find({});" }; + + const out = await h.schedule( + [ + PULL_JOB, + { name: 'bare_body', schedule: INTERVAL, body }, + { name: 'bare_fn', schedule: INTERVAL, handler: 'sweep' }, + { ...PULL_JOB, name: 'scoped_pull', organization: 'org_a' }, + ], + { functions: { sweep: async () => undefined } }, + ); + + expect(out.missingOrganization).toEqual(['orders_pull_hourly', 'bare_body', 'bare_fn']); + expect(out.pulls).toEqual(['scoped_pull']); + expect([...h.scheduled.keys()]).toEqual(['scoped_pull']); + const errors = h.said('error'); + for (const name of ['orders_pull_hourly', 'bare_body', 'bare_fn']) { + expect(errors.some((m) => m.includes('NOT SCHEDULED') && m.includes(`'${name}'`) && m.includes('`organization`'))).toBe(true); + } + }); + + it('an empty-string organization is not a declaration — it is refused like none', async () => { + const h = harness(); + const out = await h.schedule([{ ...PULL_JOB, organization: '' }]); + expect(out.missingOrganization).toEqual(['orders_pull_hourly']); + }); +}); + +describe('#20281 stage ③ (Q2-O1): under group, an undeclared job is scheduled and named once', () => { + withPosture('group'); + + it('schedules both; one warn names only the undeclared job', async () => { + const h = harness(); + const out = await h.schedule([PULL_JOB, { ...PULL_JOB, name: 'scoped_pull', organization: 'org_a' }]); + + expect(out.pulls).toEqual(['orders_pull_hourly', 'scoped_pull']); + expect(out.missingOrganization).toEqual([]); + const warns = h.said('warn').filter((m) => m.includes('NO acting organization')); + expect(warns).toHaveLength(1); + expect(warns[0]).toContain('orders_pull_hourly'); + expect(warns[0]).not.toContain('scoped_pull'); + }); +}); + +describe('#20281 stage ③: an unreadable tenancy posture fails closed — and is never read with the switch off', () => { + withPosture('isolatd'); + + it('switch on: nothing is scheduled, withheld scheduled-work-policy-unreadable, said at error', async () => { + const h = harness(); + const out = await h.schedule([PULL_JOB, { ...PULL_JOB, name: 'scoped_pull', organization: 'org_a' }]); + + expect(out.withheld).toBe('scheduled-work-policy-unreadable'); + expect(h.scheduled.size).toBe(0); + expect(h.said('error').some((m) => m.includes('scheduled-work policy could not be read'))).toBe(true); + }); + + it('switch off: the posture is not read — the deployment-policy withholding answers, as before', async () => { + const prior = process.env[SCHEDULED_WORK_ENV]; + delete process.env[SCHEDULED_WORK_ENV]; + try { + const h = harness(); + const out = await h.schedule([PULL_JOB]); + expect(out.withheld).toBe('scheduled-work-disabled'); + expect(h.said('error')).toEqual([]); + } finally { + if (prior === undefined) delete process.env[SCHEDULED_WORK_ENV]; + else process.env[SCHEDULED_WORK_ENV] = prior; + } + }); +}); diff --git a/packages/runtime/src/app-artifact-handlers.jobs.test.ts b/packages/runtime/src/app-artifact-handlers.jobs.test.ts index 577da57d6e6..3e09c18d955 100644 --- a/packages/runtime/src/app-artifact-handlers.jobs.test.ts +++ b/packages/runtime/src/app-artifact-handlers.jobs.test.ts @@ -270,9 +270,13 @@ describe('#21489: scheduleAppArtifactJobs — handler jobs and the door-wide gat const out = await h.schedule(pkg([ { name: 'off_body', schedule: INTERVAL, body: WRITE_BODY, enabled: false }, { name: 'off_handler', schedule: INTERVAL, handler: 'tick', enabled: false }, + // #20281 stage ③: the pull run form, disabled, is skipped before it is judged. + { name: 'off_pull', schedule: INTERVAL, pull: { mapping: 'nope' }, enabled: false }, ])); - expect(out).toEqual({ bodies: [], handlers: [], notScheduled: [], failed: [], cancelled: [] }); + expect(out).toEqual({ + bodies: [], handlers: [], pulls: [], notScheduled: [], missingOrganization: [], failed: [], cancelled: [], + }); expect(h.jobs.scheduled.size).toBe(0); }); diff --git a/packages/runtime/src/app-artifact-handlers.ts b/packages/runtime/src/app-artifact-handlers.ts index f518b4feaaa..99edff5912f 100644 --- a/packages/runtime/src/app-artifact-handlers.ts +++ b/packages/runtime/src/app-artifact-handlers.ts @@ -65,13 +65,35 @@ * (the job service and the engine have registered), while install-local's * doors are already past that point. One implementation, two moments. * - * A job runs on a JSON door only through a `body` that binds. Its deprecated - * `handler` names a `defineStack({ functions })` entry, which is code: it - * travels in the artifact's runtime module, which only `os start --artifact` - * loads, so no JSON door can ever resolve it. {@link collectJobsWithoutBody} - * names those jobs — and the ones whose `body` the declaration refuses - * (`judgeJobBody`) — and the install-local install route refuses a package that - * declares one enabled. + * A job runs on a JSON door only through a `body` that binds, or a `pull`. Its + * deprecated `handler` names a `defineStack({ functions })` entry, which is + * code: it travels in the artifact's runtime module, which only `os start + * --artifact` loads, so no JSON door can ever resolve it. + * {@link collectJobsWithoutBody} names those jobs — and the ones whose `body` + * the declaration refuses (`judgeJobBody`) — and the install-local install route + * refuses a package that declares one enabled. + * + * ## The pull run form, and the organization a job runs as (#20281 stage ③) + * + * `JobSchema.pull` (ruling Q1-B) is the third run form, and the declarative + * one: `{ mapping }` names a mapping whose `connectorSource` the job pulls on + * its schedule. It carries no code, so it binds HERE, in this one binder — + * {@link judgeJobPull} judges it, and each run calls the `automation` service's + * contract method, `IAutomationService.pullConnectorSource`, resolved through + * the service registry. A refused pull rejects the run (`failed`, retried per + * `retryPolicy`); a pull whose rows the import runner refused resolves + * `degraded` ({@link pullRunOutcomeOf}). It is data, like a `body`, so + * {@link collectJobsWithoutBody} never names it. + * + * `JobSchema.organization` (ruling Q2-O1) is the organization a job runs as — + * its `body`, its `handler` and its `pull` alike — judged here, at bind, by the + * deployment-posture rule scheduled flows already use + * (`resolveScheduledWorkPolicy`, `@objectstack/types`): required under + * `isolated` (a job that declares none is NOT scheduled, at `error`), optional + * under `group` (an undeclared job runs with no organization, said once at + * `warn`), not required under `single`. Each run's execution context is + * `{ isSystem: true, tenantId: }`, or `{ isSystem: true }` for a + * job that declares none ({@link jobExecutionContext}). * * A job's identity is the package's and the job's name together, as the * metadata registry keys it (`:`), so two packages may each @@ -101,8 +123,23 @@ */ import type { PluginContext } from '@objectstack/core'; -import type { IJobService, IObjectQLEngine, JobHandler, Logger } from '@objectstack/spec/contracts'; -import { resolveScheduledWorkEnabled, SCHEDULED_WORK_DISABLED_REASON } from '@objectstack/types'; +import type { + ConnectorSourcePullResult, + IAutomationService, + IJobService, + IObjectQLEngine, + JobHandler, + JobRunOutcome, + Logger, +} from '@objectstack/spec/contracts'; +import { JobSchema } from '@objectstack/spec/system'; +import { ScheduleOrganizationSchema } from '@objectstack/spec/automation'; +import { + resolveScheduledWorkEnabled, + resolveScheduledWorkPolicy, + SCHEDULED_WORK_DISABLED_REASON, + type ScheduledWorkPolicy, +} from '@objectstack/types'; import { SEMCONV } from '@objectstack/observability'; import { QuickJSScriptRunner } from './sandbox/quickjs-runner.js'; import { hookBodyRunnerFactory, actionBodyRunnerFactory, jobBodyRunnerFactory, judgeJobBody } from './sandbox/body-runner.js'; @@ -326,12 +363,19 @@ export interface JobWithoutBody { * Reads the jobs the binder reads ({@link collectBundleJobs}), and calls a job * enabled exactly when the binder does: `enabled: false` is the one value that * disables it (the schema's default is `true`). + * + * A job declaring `pull` is never named (#20281 stage ③): the pull run form is + * data, like a `body`, and binds on every door. Whether its pull binds is + * {@link judgeJobPull}'s question, answered by the binder at bind — named here + * it would read as "has no `body`" and send the author to write one beside the + * pull, which the declaration refuses. */ export function collectJobsWithoutBody(bundle: unknown): JobWithoutBody[] { const out: JobWithoutBody[] = []; for (const job of collectBundleJobs(bundle)) { if (!job || typeof job !== 'object') continue; if (job.enabled === false) continue; + if (job.pull !== undefined) continue; let bodyRefusal: string | undefined; if (job.body) { const judged = judgeJobBody(job.body); @@ -347,6 +391,141 @@ export function collectJobsWithoutBody(bundle: unknown): JobWithoutBody[] { return out; } +// ─── The pull run form, and the organization a job runs as (#20281 stage ③) ─ + +/** What {@link judgeJobPull} answers for a job that declares `pull`. */ +export type JobPullJudgement = + | { binds: true; mapping: string } + | { binds: false; refusal: string }; + +/** + * The mappings an artifact declares — the resolved top-level `mappings` + * (ADR-0130 D4: every package body's too), else the legacy `manifest.mappings`, + * the same one-list-or-the-other read {@link collectBundleJobs} makes. A + * collection in the record form (`{ : mapping }`) is read by its keys. + */ +function collectBundleMappings(bundle: unknown): Array> { + const stack = resolveArtifactCollections(bundle) as any; + const raw = stack?.mappings !== undefined ? stack.mappings : (bundle as any)?.manifest?.mappings; + if (Array.isArray(raw)) return raw.filter((m): m is Record => !!m && typeof m === 'object'); + if (raw && typeof raw === 'object') { + return Object.entries(raw as Record) + .filter(([, m]) => !!m && typeof m === 'object') + .map(([key, m]) => ({ name: key, ...(m as Record) })); + } + return []; +} + +/** + * Does a job's `pull` bind on this artifact? The ONE judgement, read by the + * binder ({@link scheduleAppArtifactJobs}) before it schedules a pull job. + * + * It binds when the job declares no code beside it (the declaration refuses + * `pull` + `body` / `handler`), `pull` parses against `JobSchema.pull`, and the + * artifact declares the named mapping WITH a `connectorSource` — the reference + * `defineStack` (and so `os validate`) checks against the stack's mappings. A + * pull that does not bind is not scheduled: a run of it could only ever be + * refused by the automation service (`mapping_not_found` / + * `no_connector_source`), so the reason is said once, at bind, rather than on + * every tick. + * + * `refusal` is the sentence the author acts on, prefixed with the key it names. + */ +export function judgeJobPull( + job: { pull?: unknown; body?: unknown; handler?: unknown }, + bundle: unknown, +): JobPullJudgement { + if (job.body !== undefined || job.handler !== undefined) { + return { + binds: false, + refusal: 'pull: the job declares `body` or `handler` beside `pull`, which the declaration refuses — the platform ' + + 'binds the pull itself, so one of the two run forms would never run; keep `pull` or keep the code', + }; + } + const parsed = JobSchema.shape.pull.safeParse(job.pull); + if (!parsed.success || parsed.data === undefined) { + const first = parsed.success ? undefined : parsed.error.issues[0]; + const at = first && first.path.length > 0 ? `pull.${first.path.map(String).join('.')}: ` : 'pull: '; + return { binds: false, refusal: `${at}${first ? first.message : 'expected `{ mapping: "" }`'}` }; + } + const mapping = parsed.data.mapping; + const declared = collectBundleMappings(bundle).find((m) => m.name === mapping); + if (!declared) { + return { + binds: false, + refusal: `pull.mapping: this artifact declares no mapping '${mapping}' — declare it (its targetObject, fieldMapping ` + + 'and a connectorSource naming the rest or openapi connector it pulls from) beside the job, or correct the name', + }; + } + if (declared.connectorSource === undefined) { + return { + binds: false, + refusal: `pull.mapping: mapping '${mapping}' declares no connectorSource, so there is nothing to pull — add ` + + 'connectorSource: { connector, action } naming the rest or openapi connector it reads from', + }; + } + return { binds: true, mapping }; +} + +/** + * The organization a job declares (`JobSchema.organization`), or `undefined`. + * A present-but-unusable value (an empty string, a number) answers `undefined` + * — the value shape is the scheduled flow's (`ScheduleOrganizationSchema`), and + * so is the reading: reporting "declared" for a value nothing can act on is the + * silent acceptance the key exists to end (`resolveScheduleOrganization`). + */ +export function resolveJobOrganization(job: { organization?: unknown }): string | undefined { + const parsed = ScheduleOrganizationSchema.safeParse(job.organization); + return parsed.success ? parsed.data : undefined; +} + +/** + * The execution context a job RUNS AS — its `body`'s `ctx.api`, its + * `handler`'s `JobHandlerContext.executionContext` and its `pull`'s target read + * and writes alike: system, carrying the job's organization as `tenantId` (the + * field the tenancy guard reads) where it declares one. Fresh per run, never a + * shared constant — an execution envelope is a value the engine may extend. + */ +export function jobExecutionContext(organization: string | undefined): { isSystem: true; tenantId?: string } { + return organization !== undefined ? { isSystem: true, tenantId: organization } : { isSystem: true }; +} + +/** + * The run outcome of a pull that RESOLVED, mapped once here rather than left to + * each author: rows the import runner refused make the run `degraded` — it ran + * to completion and part (or all) of its work did not happen, and a retry + * would refuse the same rows — with the counts as the reason; otherwise the + * run `completed`, including a pull that carried no records. A pull that was + * REFUSED never reaches here: the service rejected, which is the job's + * `failed` (and the retry policy's trigger). + */ +export function pullRunOutcomeOf(result: ConnectorSourcePullResult): JobRunOutcome { + const summary = result?.summary; + const errors = Number(summary?.errors ?? 0); + if (errors > 0) { + return { + outcome: 'degraded', + reason: `mapping '${result.mapping}': ${errors} of ${result.pulled} pulled record(s) were refused by the import runner ` + + `(${Number(summary?.created ?? 0)} created, ${Number(summary?.updated ?? 0)} updated, ${Number(summary?.skipped ?? 0)} skipped)`, + }; + } + return { outcome: 'completed' }; +} + +/** + * The sentence an `isolated` deployment's binder logs (at `error`) for a job + * that declares no organization — the job's counterpart of the scheduled + * flow's `describeMissingScheduleOrganization`: the same rule, said about a + * job's key rather than a start node's. + */ +function describeMissingJobOrganization(jobName: string): string { + return `job '${jobName}' declares no \`organization\`, and this deployment runs the 'isolated' tenancy posture ` + + 'with package-authored scheduled work switched on: a scheduled run has no session to inherit an organization ' + + 'from, so its tenant-scoped writes would all be refused — the job is NOT scheduled. Declare the organization it ' + + "runs in: `organization: ''` on the job. Work wanted in several organizations is one job " + + 'per organization — a job is never fanned out across them, and no organization is ever chosen for it.'; +} + // ─── Hooks with no `body` (#21585) ───────────────────────────────────── /** @@ -415,16 +594,26 @@ export interface AppArtifactJobSchedulingOptions { export interface AppArtifactJobScheduling { /** * Set when nothing was scheduled for a reason that holds for every job: - * the deployment does not run package-authored scheduled work (#17396), or - * no job service is registered. + * the deployment does not run package-authored scheduled work (#17396), no + * job service is registered, or (#20281 stage ③) the deployment's + * scheduled-work policy could not be read — an unrecognized tenancy + * posture, so whether a job must declare its organization is unknown. */ - withheld?: 'scheduled-work-disabled' | 'no-job-service'; + withheld?: 'scheduled-work-disabled' | 'no-job-service' | 'scheduled-work-policy-unreadable'; /** Jobs scheduled to run their sandboxed `body`. */ bodies: string[]; /** Jobs scheduled to run the `functions` entry their `handler` names. */ handlers: string[]; - /** Enabled jobs with nothing this door could run (an unbindable body, an unresolvable handler, no name). */ + /** [#20281 stage ③] Jobs scheduled to pull the mapping their `pull` names. */ + pulls: string[]; + /** Enabled jobs with nothing this door could run (an unbindable body or pull, an unresolvable handler, no name). */ notScheduled: string[]; + /** + * [#20281 stage ③] Enabled jobs NOT scheduled because they declare no + * `organization` on a deployment whose posture requires one (`isolated`, + * with scheduled work switched on). + */ + missingOrganization: string[]; /** Jobs whose `IJobService.schedule` call threw. */ failed: string[]; /** @@ -445,6 +634,14 @@ export interface AppArtifactJobScheduling { * Per job, in this order: * * - `enabled: false` → skipped (debug); + * - (#20281 stage ③) no `organization` on a deployment whose scheduled-work + * policy requires one (`isolated`, switch on) → NOT scheduled, at `error` + * ({@link describeMissingJobOrganization}), whatever its run form; + * - a `pull` → the mapping pull, when {@link judgeJobPull} binds it and the + * `automation` service serves `pullConnectorSource`; each run calls that + * contract method through the service registry under the job's execution + * context, and maps the result with {@link pullRunOutcomeOf}. A pull that + * does not bind schedules nothing (warn); * - a `body` → the sandboxed body (`jobBodyRunnerFactory`), and the `body` * WINS when a `handler` is declared beside it. A body that cannot be bound * (wrong shape, a `body.timeoutMs`) schedules nothing — never the handler @@ -455,6 +652,11 @@ export interface AppArtifactJobScheduling { * door refuses the shape up front ({@link collectJobsWithoutBody}); * - else → not scheduled (warn). * + * Every form runs as {@link jobExecutionContext} of the job's organization — + * the body's `ctx.api`, the handler's `executionContext`, the pull's + * `context`. Under `group`, the jobs that declare none are named once at + * `warn`: they are armed, and a tenant-scoped write they make is refused. + * * The schedule is lowered to the boundary tier (`toBoundaryJobSchedule`), and * the job's `retryPolicy` / `timeoutMs` are threaded to the adapter. For a body * job the same `timeoutMs` also bounds the sandbox run — the one limit @@ -488,7 +690,9 @@ export async function scheduleAppArtifactJobs( const { appId, ql } = options; const logger: Logger = ctx.logger; const tag = `[${options.source ?? 'AppPlugin'}]`; - const out: AppArtifactJobScheduling = { bodies: [], handlers: [], notScheduled: [], failed: [], cancelled: [] }; + const out: AppArtifactJobScheduling = { + bodies: [], handlers: [], pulls: [], notScheduled: [], missingOrganization: [], failed: [], cancelled: [], + }; const jobs = collectBundleJobs(bundle); let svc: IJobService | undefined; @@ -529,6 +733,28 @@ export async function scheduleAppArtifactJobs( logger.warn(`${tag} job service not registered — skipping declarative jobs`, { appId, jobCount: jobs.length }); return { ...out, withheld: 'no-job-service' }; } + // [#20281 stage ③] The posture half of the scheduled-work policy — read + // once per call, AFTER the switch: whether a job must declare its + // organization is the rule scheduled flows bind by + // (`requiresActingOrganization`, `runOwnership`), and it is only asked of + // a deployment that runs scheduled work at all. An unrecognized tenancy + // posture makes the resolver throw (a typo must not resolve to `single` + // and drop the requirement with it), so this fails CLOSED: nothing is + // scheduled, and the reason is said at `error`. + let policy: ScheduledWorkPolicy; + try { + policy = resolveScheduledWorkPolicy(); + } catch (err: any) { + logger.error( + `${tag} declarative jobs NOT scheduled — the deployment's scheduled-work policy could not be read, so whether ` + + 'each job must declare its `organization` is unknown. Correct the tenancy posture the error names, then restart.', + err as Error, + { appId, jobCount: jobs.length }, + ); + out.cancelled = await retireAppJobs(svc, appId, new Set(), logger, tag); + return { ...out, withheld: 'scheduled-work-policy-unreadable' }; + } + const jobService: IJobService = svc; ensureJobUninstallCleanup(ctx, jobService); @@ -541,6 +767,10 @@ export async function scheduleAppArtifactJobs( const collections = resolveArtifactCollections(bundle); const bodyRunner = jobBodyRunnerFactory(new QuickJSScriptRunner(), { ql, logger, appId }); const metrics = resolveMetrics(ctx); + // [#20281 stage ③] Jobs armed with NO organization on a posture whose + // writes need one (`group`): legal, and refused at their first + // tenant-scoped write — said once below, with the remedy. + const armedWithoutOrganization: string[] = []; for (const job of jobs) { const jobName: string = job?.name; @@ -554,12 +784,72 @@ export async function scheduleAppArtifactJobs( continue; } + // [#20281 stage ③] The organization the job runs as, judged by the + // scheduled flows' posture rule. ⛔ No limb picks one for a job that + // declares none: under `isolated` it is not scheduled, and elsewhere + // it runs with none — a wrong `organization_id` is worse than a + // refusal, because it is silently authoritative to every reader. + const organization = resolveJobOrganization(job); + if (organization === undefined && policy.requiresActingOrganization) { + logger.error(`${tag} NOT SCHEDULED — ${describeMissingJobOrganization(jobName)}`, undefined, { + appId, + job: jobName, + posture: policy.posture, + }); + out.missingOrganization.push(jobName); + continue; + } + if (organization === undefined && policy.runOwnership === 'per-record') { + armedWithoutOrganization.push(jobName); + } + let run: JobHandler; - let form: 'body' | 'handler'; - if (job.body) { + let form: 'body' | 'handler' | 'pull'; + if (job.pull !== undefined) { + // The declarative run form. Judged against the artifact first: a + // pull whose mapping is missing (or has no `connectorSource`) + // could only ever be refused, so it is not scheduled. + const judged = judgeJobPull(job, bundle); + if (!judged.binds) { + logger.warn(`${tag} job pull does not bind — the job is NOT scheduled: ${judged.refusal}`, { + appId, + job: jobName, + }); + out.notScheduled.push(jobName); + continue; + } + const mapping = judged.mapping; + let automation: IAutomationService | undefined; + try { automation = ctx.getService('automation'); } catch { /* not installed */ } + if (typeof automation?.pullConnectorSource !== 'function') { + logger.warn( + `${tag} job '${jobName}' pulls mapping '${mapping}', but no \`automation\` service serving ` + + '`pullConnectorSource` is registered — the job is NOT scheduled. Compose AutomationServicePlugin ' + + '(`@objectstack/service-automation`), which holds the connector registry a pull reads through.', + { appId, job: jobName, mapping }, + ); + out.notScheduled.push(jobName); + continue; + } + // Resolved again on every run, through the registry — never the + // instance read above, so a re-registered service is the one called. + run = async () => { + let service: IAutomationService | undefined; + try { service = ctx.getService('automation'); } catch { /* reported below */ } + if (typeof service?.pullConnectorSource !== 'function') { + throw new Error( + `job '${jobName}' pulls mapping '${mapping}', but the \`automation\` service no longer serves ` + + '`pullConnectorSource` — nothing was pulled', + ); + } + const result = await service.pullConnectorSource({ mapping, context: jobExecutionContext(organization) }); + return pullRunOutcomeOf(result); + }; + form = 'pull'; + } else if (job.body) { // The body wins over a `handler` beside it, as for hooks. When it // cannot be bound the factory has said why, and nothing runs. - const bound = bodyRunner(job); + const bound = bodyRunner({ ...job, organization }); if (!bound) { out.notScheduled.push(jobName); continue; @@ -593,6 +883,9 @@ export async function scheduleAppArtifactJobs( bundle: collections, ql: ql as IObjectQLEngine, logger, + // [#20281 stage ③] The envelope this job runs as — the one + // its `body` or `pull` would carry. `ql` stays the raw engine. + executionContext: jobExecutionContext(organization), }; // #14256: RETURN the handler's resolved value. `JobHandler` is // `(context) => Promise` and all three @@ -622,7 +915,7 @@ export async function scheduleAppArtifactJobs( ? { retryPolicy: job.retryPolicy, timeoutMs: job.timeoutMs } : undefined, ); - (form === 'body' ? out.bodies : out.handlers).push(jobName); + (form === 'body' ? out.bodies : form === 'pull' ? out.pulls : out.handlers).push(jobName); claimJobKey(jobService, appId, jobName, key); if (heldBy !== undefined) { logger.info( @@ -648,14 +941,38 @@ export async function scheduleAppArtifactJobs( } } - out.cancelled = await retireAppJobs(jobService, appId, new Set([...out.bodies, ...out.handlers]), logger, tag); + out.cancelled = await retireAppJobs( + jobService, + appId, + new Set([...out.bodies, ...out.handlers, ...out.pulls]), + logger, + tag, + ); + + // [#20281 stage ③] Said once, at `warn`, the way the scheduled flows' + // bind line says it for a record-less flow under `group`: armed and legal, + // and its first tenant-scoped write is refused — boot is where an operator + // is reading, so the remedy is given here rather than only at that tick. + const armed = new Set([...out.bodies, ...out.handlers, ...out.pulls]); + const unscoped = armedWithoutOrganization.filter((name) => armed.has(name)); + if (unscoped.length > 0) { + logger.warn( + `${tag} ${unscoped.length} job(s) scheduled with NO acting organization (tenancy posture '${policy.posture}'): ` + + `${unscoped.join(', ')}. A job has no record to derive one from, so any tenant-scoped row it writes is REFUSED ` + + 'at the write. Declare `organization` on the job if it writes per-organization data.', + { appId, jobs: unscoped, posture: policy.posture }, + ); + } - const scheduled = out.bodies.length + out.handlers.length; + const scheduled = out.bodies.length + out.handlers.length + out.pulls.length; logger.info(`${tag} Scheduled background jobs`, { appId, count: scheduled, bodies: out.bodies.length, handlers: out.handlers.length, + pulls: out.pulls.length, + notScheduled: out.notScheduled.length, + missingOrganization: out.missingOrganization.length, failed: out.failed.length, cancelled: out.cancelled.length, }); diff --git a/packages/runtime/src/app-plugin.job-data-reach.test.ts b/packages/runtime/src/app-plugin.job-data-reach.test.ts index b6ee6257ec4..0e7574cd8ae 100644 --- a/packages/runtime/src/app-plugin.job-data-reach.test.ts +++ b/packages/runtime/src/app-plugin.job-data-reach.test.ts @@ -92,6 +92,12 @@ const NOTE = { const PRE_14094_KEYS = ['bundle', 'data', 'jobId'] as const; /** What #14094 added, and nothing else. */ const ADDED_KEYS = ['logger', 'ql'] as const; +/** + * What #20281 stage ③ added (ruling Q2-O1): the envelope the job RUNS AS — + * `{ isSystem: true, tenantId }` from the job's declared `organization`, else + * `{ isSystem: true }`. Additive, like #14094's: `ql` stays the raw engine. + */ +const ORGANIZATION_KEYS = ['executionContext'] as const; interface Harness { engine: ObjectQL; @@ -244,7 +250,7 @@ describe('#14094 — a declarative job handler has data reach (TS-config path)', expect(h.errorLogs()).toEqual([]); }); - it('the context is the pre-#14094 set PLUS exactly `ql` and `logger`', async () => { + it('the context is the pre-#14094 set PLUS exactly `ql` and `logger` — and #20281\'s `executionContext`', async () => { // Reads nothing and writes nothing — this one is about the shape. const h = await harness(); const seen: Array> = []; @@ -260,7 +266,7 @@ describe('#14094 — a declarative job handler has data reach (TS-config path)', expect(seen).toHaveLength(1); const keys = Object.keys(seen[0]).sort(); - expect(keys).toEqual([...PRE_14094_KEYS, ...ADDED_KEYS].sort()); + expect(keys).toEqual([...PRE_14094_KEYS, ...ADDED_KEYS, ...ORGANIZATION_KEYS].sort()); // The pre-existing members keep their meaning — `jobId` is the job's // name, `bundle` is the metadata bundle, `data` is the trigger payload. @@ -271,6 +277,8 @@ describe('#14094 — a declarative job handler has data reach (TS-config path)', // And the added members are the LIVE handles, not placeholders. expect(seen[0].ql).toBe(h.engine); expect(seen[0].logger).toBe(h.ctx.logger); + // A job that declares no organization runs as plain system. + expect(seen[0].executionContext).toEqual({ isSystem: true }); }); it('`data` from a manual trigger still reaches the handler beside the new members', async () => { diff --git a/packages/runtime/src/job-handler-context.ts b/packages/runtime/src/job-handler-context.ts index d9d964f3feb..4274d17d201 100644 --- a/packages/runtime/src/job-handler-context.ts +++ b/packages/runtime/src/job-handler-context.ts @@ -91,4 +91,24 @@ export interface JobHandlerContext { * diagnostics land in the platform's log stream instead of `console`. */ logger: Logger; + /** + * The execution context this job RUNS AS — the envelope the binder builds + * from the job's declared `organization` (`JobSchema.organization`, judged + * at bind by the scheduled-work posture rule; ruling Q2-O1 on the + * connector-sync card): `{ isSystem: true, tenantId: '' }`, + * or `{ isSystem: true }` for a job that declares none. The same envelope a + * `body` job's `ctx.api` and a `pull` job's reads and writes carry. + * + * `ql` is the raw engine and stays so (an existing handler is unchanged byte + * for byte), so a handler writes as the job's organization by passing this + * as each call's `context` — `ql.insert('task', row, { context: + * executionContext })`. A call that passes none carries no organization, + * and under the `group` / `isolated` postures a tenant-scoped system write + * without one is refused at the write. + * + * Always set when the binder invokes the handler; optional in the type only + * so that code which BUILDS a context (a test calling a handler directly) + * keeps compiling — the widening stays additive in both directions. + */ + executionContext?: { isSystem: true; tenantId?: string }; } diff --git a/packages/runtime/src/sandbox/body-runner.ts b/packages/runtime/src/sandbox/body-runner.ts index 4a2653b2c52..4820df85a92 100644 --- a/packages/runtime/src/sandbox/body-runner.ts +++ b/packages/runtime/src/sandbox/body-runner.ts @@ -536,6 +536,14 @@ export function judgeJobBody(raw: unknown): JobBodyJudgement { * declared `capabilities` and the stored-metadata boundary every body's api * carries ({@link buildSandboxApi}). * + * …as SYSTEM IN the job's organization, when the binder hands one over + * (`job.organization` — `JobSchema.organization`, judged at bind by the + * scheduled-work posture rule): the envelope is `{ isSystem: true, tenantId }`, + * so a tenant-scoped write carries that organization the way a session write + * does, and is no longer refused under the `group` / `isolated` postures for + * want of one (ruling Q2-O1 on the connector-sync card). With none it stays + * `{ isSystem: true }`. + * * ## The time limit * * The job's own `timeoutMs` reaches the runner as `opts.timeoutMs` — the ONE @@ -560,7 +568,7 @@ export function judgeJobBody(raw: unknown): JobBodyJudgement { export function jobBodyRunnerFactory( runner: ScriptRunner, opts: FactoryOptions, -): (job: { name: string; body?: unknown; timeoutMs?: number }) => JobHandler | undefined { +): (job: { name: string; body?: unknown; timeoutMs?: number; organization?: string }) => JobHandler | undefined { return (job) => { const raw = job.body; if (!raw) return undefined; @@ -580,6 +588,7 @@ export function jobBodyRunnerFactory( const sandboxCtx = buildJobSandboxContext( opts.ql, buildBodyLogSurface(opts, { kind: 'job', name: job.name }), + job.organization, ); try { opts.logger?.debug?.('[BodyRunner] job fired', { appId: opts.appId, job: job.name }); @@ -1180,11 +1189,17 @@ function buildActionSandboxContext( * refuses an owner-scoped write that has neither a `userId` to own it nor * `isSystem` to bypass. Fresh per run, never a shared constant, because an * execution envelope is a value the engine may extend (a transaction joins it). + * + * `organization` — the job's declared one, which the binder resolved — joins + * the envelope as `tenantId`, the field the tenancy guard reads + * (`resolveSystemWriteOrganization`'s remedy: "pass it on the execution + * context"). Absent, no `tenantId` key is written at all. */ -function buildJobSandboxContext(ql: any, log: ScriptContext['log']): ScriptContext { +function buildJobSandboxContext(ql: any, log: ScriptContext['log'], organization?: string): ScriptContext { + const executionContext = organization ? { isSystem: true, tenantId: organization } : { isSystem: true }; return { input: undefined, - api: buildSandboxApi({ executionContext: { isSystem: true } }, ql, 'job body'), + api: buildSandboxApi({ executionContext }, ql, 'job body'), log, crypto: globalThis.crypto, }; diff --git a/packages/services/service-automation/src/connector-pull-service-door.test.ts b/packages/services/service-automation/src/connector-pull-service-door.test.ts new file mode 100644 index 00000000000..b951d2b436c --- /dev/null +++ b/packages/services/service-automation/src/connector-pull-service-door.test.ts @@ -0,0 +1,76 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +/** + * #20281 stage ③ — the connector sync executor on the `automation` SERVICE. + * + * Before this, `pullConnectorSource` was reachable only as an + * `AutomationServicePlugin` instance method, which no other package holds: the + * registered `automation` service is the ENGINE, and the contract + * (`IAutomationService`) declared no pull. A job's `pull` run form binds + * through the service registry, so the contract gains `pullConnectorSource` + * and the engine serves it from the executor the plugin attaches at `init()` + * (`setConnectorPullSource`) — the plugin keeps the materialized-connector map, + * the engine is handed the call. + * + * Pinned: the bare engine refuses with an ADR-0112 envelope rather than + * answering a pull that never ran; an attached executor receives the request + * as sent and its answer comes back unchanged; and a booted kernel's + * `automation` service reaches the plugin's executor. The end-to-end pull + * through the service, over a real `rest` connector, is in + * `connector-pull.integration.test.ts`. + */ + +import { describe, expect, it, vi } from 'vitest'; +import { LiteKernel } from '@objectstack/core'; +import type { ConnectorSourcePullResult, IAutomationService } from '@objectstack/spec/contracts'; +import { AutomationEngine } from './engine.js'; +import { AutomationServicePlugin } from './plugin.js'; + +const silent = { debug: () => {}, info: () => {}, warn: () => {}, error: () => {} } as any; + +const RESULT: ConnectorSourcePullResult = { + mapping: 'orders_pull', + targetObject: 'order', + connector: 'orders_api', + action: 'request', + pulled: 1, + summary: { total: 1, processed: 1, created: 1, updated: 0, skipped: 0, errors: 0, ok: 1, cancelled: false }, +}; + +describe('#20281 stage ③: IAutomationService.pullConnectorSource, served by the engine', () => { + it('a bare engine — no executor attached — refuses with SERVICE_UNAVAILABLE 503 and pulls nothing', async () => { + const engine = new AutomationEngine(silent); + await expect(engine.pullConnectorSource({ mapping: 'orders_pull' })).rejects.toMatchObject({ + code: 'SERVICE_UNAVAILABLE', + status: 503, + }); + }); + + it('an attached executor receives the request as sent, and its answer comes back unchanged', async () => { + const engine = new AutomationEngine(silent); + const source = vi.fn(async () => RESULT); + engine.setConnectorPullSource(source); + + const request = { mapping: 'orders_pull', context: { isSystem: true, tenantId: 'org_a' } }; + await expect(engine.pullConnectorSource(request)).resolves.toBe(RESULT); + expect(source).toHaveBeenCalledWith(request); + }); + + it('a booted kernel\'s `automation` service reaches the plugin\'s executor — the door a job pull binds through', async () => { + const plugin = new AutomationServicePlugin({ suspendedRunStore: 'memory' }); + const executor = vi.spyOn(plugin, 'pullConnectorSource').mockResolvedValue(RESULT as any); + const kernel = new LiteKernel({ logger: { level: 'silent' } } as never); + kernel.use(plugin); + await kernel.bootstrap(); + try { + const service = kernel.getService('automation'); + expect(typeof service.pullConnectorSource).toBe('function'); + + const request = { mapping: 'orders_pull', context: { isSystem: true } }; + await expect(service.pullConnectorSource!(request)).resolves.toBe(RESULT); + expect(executor).toHaveBeenCalledWith(request); + } finally { + await kernel.shutdown(); + } + }); +}); diff --git a/packages/services/service-automation/src/connector-pull.integration.test.ts b/packages/services/service-automation/src/connector-pull.integration.test.ts index ba8265e4662..01bad9c8260 100644 --- a/packages/services/service-automation/src/connector-pull.integration.test.ts +++ b/packages/services/service-automation/src/connector-pull.integration.test.ts @@ -32,6 +32,7 @@ import { SqlDriver } from '@objectstack/driver-sql'; import { ConnectorRestPlugin } from '../../../connectors/connector-rest/src/index.js'; import { AutomationServicePlugin } from './plugin.js'; import type { EngineQueryOptions } from '@objectstack/spec/data'; +import type { IAutomationService } from '@objectstack/spec/contracts'; import type { ConnectorPullProtocol } from './connector-pull.js'; const CONTACT = { @@ -180,10 +181,17 @@ describe('[#20919] connector pull, end to end through a real rest connector', () expect(first.watermark).toEqual({ field: 'updated_at', target: 'synced_at', param: 'since', from: undefined }); // Pull 2: the starting point is the highest `synced_at` stored by pull 1. - const second = await automation.pullConnectorSource({ mapping: 'crm_contacts', context: SYSTEM }); + // [#20281 stage ③] Taken through the kernel's `automation` SERVICE — + // `IAutomationService.pullConnectorSource`, the door a job's `pull` run + // form binds through — not the plugin instance: the engine serves it + // from the executor the plugin attached at init(). + const service = kernel.getService('automation'); + const second = await service.pullConnectorSource!({ mapping: 'crm_contacts', context: SYSTEM }); expect(fixture.seen[1].get('since')).toBe('2026-01-02T00:00:00.000Z'); expect(second.watermark?.from).toBe('2026-01-02T00:00:00.000Z'); - expect(second.summary.errors, JSON.stringify(second.summary.results)).toBe(0); + // The contract declares the runner's tallies, not its per-row results, + // so the failure message carries the whole summary as the service sent it. + expect(second.summary.errors, JSON.stringify(second.summary)).toBe(0); expect({ pulled: second.pulled, created: second.summary.created, updated: second.summary.updated }) .toEqual({ pulled: 2, created: 1, updated: 1 }); diff --git a/packages/services/service-automation/src/connector-pull.ts b/packages/services/service-automation/src/connector-pull.ts index 3886fca84e1..6208562d93c 100644 --- a/packages/services/service-automation/src/connector-pull.ts +++ b/packages/services/service-automation/src/connector-pull.ts @@ -27,8 +27,12 @@ * limit is stated where authors read it (`mapping.zod.ts`'s `watermark` * describe and `SYNC_ARCHITECTURE.md`). * - * ⛔ Nothing here schedules a pull. A `job` drives it (stage ③), and the - * caller hands in the execution context the pull reads and writes under. + * ⛔ Nothing here schedules a pull. A `job` whose `pull` names the mapping + * drives it (`JobSchema.pull`): the runtime's job binder calls the `automation` + * service's `pullConnectorSource` (the `IAutomationService` contract method) + * on each run, and hands in the execution context the pull reads and writes + * under — `{ isSystem: true, tenantId }` from the job's declared + * `organization`, or `{ isSystem: true }` where it declares none. * * Every refusal is loud and typed ({@link ConnectorPullError}): a pull that * cannot honour its binding throws before anything is written, and the @@ -127,7 +131,7 @@ export interface ConnectorPullDeps { export interface ConnectorPullOptions { /** Name of the `mapping` whose `connectorSource` is pulled. */ mapping: string; - /** Execution context the target read and the writes run under — the caller's (a `job`, from stage ③). */ + /** Execution context the target read and the writes run under — the caller's (a `job`'s, built from its `organization`). */ context?: any; environmentId?: string; /** Automation context handed to the connector action, as a flow node would hand it. */ diff --git a/packages/services/service-automation/src/engine.ts b/packages/services/service-automation/src/engine.ts index 7372f2a8f4f..815bfbdaddf 100644 --- a/packages/services/service-automation/src/engine.ts +++ b/packages/services/service-automation/src/engine.ts @@ -10,7 +10,7 @@ import type { FlowFunctionEffect, FlowRunSummary, } from '@objectstack/spec/automation'; -import type { AutomationContext, AutomationResult, ResumeSignal, IAutomationService, RunListResult, ScreenSpec, ScreenFieldSpec } from '@objectstack/spec/contracts'; +import type { AutomationContext, AutomationResult, ResumeSignal, IAutomationService, RunListResult, ScreenSpec, ScreenFieldSpec, ConnectorSourcePullRequest, ConnectorSourcePullResult } from '@objectstack/spec/contracts'; import { RESUME_AUTHORITY_SERVICE } from '@objectstack/spec/contracts'; import { validateScreenInputs, @@ -2416,6 +2416,8 @@ export class AutomationEngine implements IAutomationService { private packagedFlowSource?: PackagedFlowSource; /** [#20790] The write-only flow credential channel — see {@link setFlowCredentialSource}. */ private flowCredentialSource?: FlowCredentialSource; + /** [#20281 stage ③] The connector sync executor — see {@link setConnectorPullSource}. */ + private connectorPullSource?: (request: ConnectorSourcePullRequest) => Promise; /** * Re-entrancy guard for record-triggered flows (complements the intra-run * {@link MAX_NODE_REENTRIES} back-edge guard, which cannot see a self-trigger @@ -4754,6 +4756,44 @@ export class AutomationEngine implements IAutomationService { this.flowCredentialSource = source; } + /** + * [#20281 stage ③] Attach the connector sync executor this engine serves as + * {@link pullConnectorSource}. The automation plugin calls this at + * `init()`, before it registers the engine as the `automation` service: + * the executor resolves connectors against the instances the PLUGIN + * materialized, which the engine does not hold, so the plugin hands the + * engine the call rather than the map. + */ + setConnectorPullSource( + source: ((request: ConnectorSourcePullRequest) => Promise) | undefined, + ): void { + this.connectorPullSource = source; + } + + /** + * [#20281 stage ③] `IAutomationService.pullConnectorSource` — pull one + * `mapping`'s `connectorSource` and write the records through the import + * runner, under `request.context`. The door a job's `pull` run form binds + * through. Rejects, naming the gap, on an engine no plugin attached an + * executor to: there is nothing that could pull, and answering as if a + * pull ran would be the silent no-op the contract's rejection rule exists + * to rule out. + */ + async pullConnectorSource(request: ConnectorSourcePullRequest): Promise { + if (!this.connectorPullSource) { + // ADR-0112: the code the executor itself answers for "cannot pull + // now" (a degraded connector), so a caller branches on one pair. + throw Object.assign( + new Error( + `[Automation] pullConnectorSource('${request.mapping}'): this automation engine has no connector sync ` + + 'executor attached — it is attached by AutomationServicePlugin at init(); a bare engine cannot pull. Nothing was pulled.', + ), + { code: 'SERVICE_UNAVAILABLE', status: 503 }, + ); + } + return this.connectorPullSource(request); + } + /** [#20790] Does the credential channel hold the LIVE credential at this position? */ holdsFlowCredential(flowName: string, nodeId: string, key: string): boolean { return this.flowCredentialSource?.holds(flowName, nodeId, key) ?? false; diff --git a/packages/services/service-automation/src/index.ts b/packages/services/service-automation/src/index.ts index 1fa7bdbbbd8..9add96a9274 100644 --- a/packages/services/service-automation/src/index.ts +++ b/packages/services/service-automation/src/index.ts @@ -167,8 +167,10 @@ export { AutomationServicePlugin, createPackageFileLoader } from './plugin.js'; export type { AutomationServicePluginOptions } from './plugin.js'; // [#20919] The connector sync executor — pull a `mapping`'s `connectorSource` -// and write the records through the import runner. Nothing schedules it: a -// `job` drives a pull (`AutomationServicePlugin.pullConnectorSource`). +// and write the records through the import runner. A `job` whose `pull` names +// the mapping drives it, through the `automation` service's contract method +// (`IAutomationService.pullConnectorSource`, which the engine serves from +// `AutomationServicePlugin.pullConnectorSource`). export { pullConnectorSource, ConnectorPullError, CONNECTOR_PULL_PROVIDERS } from './connector-pull.js'; export type { ConnectorPullDeps, diff --git a/packages/services/service-automation/src/plugin.ts b/packages/services/service-automation/src/plugin.ts index d6d14f0d429..85be0fe6fdb 100644 --- a/packages/services/service-automation/src/plugin.ts +++ b/packages/services/service-automation/src/plugin.ts @@ -636,8 +636,13 @@ export class AutomationServicePlugin implements Plugin { * the `protocol` service, and resolves the connector against the instances * this plugin materialized from `connectors[]`. * - * ⛔ Nothing calls this on a schedule: a `job` drives a pull, and the - * caller supplies the execution context (`opts.context`) it runs under. + * The engine serves this as the `automation` service's contract method + * (`IAutomationService.pullConnectorSource`): `init()` attaches it with + * `setConnectorPullSource`, because the registered service is the ENGINE + * and the materialized-connector map it needs is this plugin's. A `job` + * whose `pull` names the mapping reaches it that way, through the service + * registry, and supplies the execution context (`opts.context`) built from + * the job's `organization`. */ async pullConnectorSource(opts: ConnectorPullOptions): Promise { const ctx = this.ctx; @@ -856,6 +861,12 @@ export class AutomationServicePlugin implements Plugin { this.credentialChannel = new FlowCredentialChannel(() => this.resolveDataEngine(ctx)); this.engine.setFlowCredentialSource(this.credentialChannel); + // [#20281 stage ③] The connector sync executor, served on the + // `automation` service's contract (`pullConnectorSource`) — the door a + // job's `pull` run form binds through. Attached before the service is + // registered, so no caller can resolve the service without it. + this.engine.setConnectorPullSource((request) => this.pullConnectorSource(request)); + // Register as global service — other plugins access via ctx.getService('automation') ctx.registerService('automation', this.engine); diff --git a/packages/spec/api-surface/contracts.json b/packages/spec/api-surface/contracts.json index 53e26b0ac1a..85f72a9b02d 100644 --- a/packages/spec/api-surface/contracts.json +++ b/packages/spec/api-surface/contracts.json @@ -61,6 +61,9 @@ "CheckNamespaceInput (interface)", "CheckNamespaceResult (interface)", "ClusterCallContext (interface)", + "ConnectorSourcePullRequest (interface)", + "ConnectorSourcePullResult (interface)", + "ConnectorSourcePullSummary (interface)", "CoreServiceContract (type)", "CoreServiceContracts (interface)", "CounterIncrOptions (interface)", diff --git a/packages/spec/authorable-surface/system.json b/packages/spec/authorable-surface/system.json index 77699514df3..bb139fbe61a 100644 --- a/packages/spec/authorable-surface/system.json +++ b/packages/spec/authorable-surface/system.json @@ -481,6 +481,8 @@ "system/Job:handler", "system/Job:label", "system/Job:name", + "system/Job:organization", + "system/Job:pull", "system/Job:retryPolicy", "system/Job:schedule", "system/Job:timeout [RETIRED]", diff --git a/packages/spec/docs/SYNC_ARCHITECTURE.md b/packages/spec/docs/SYNC_ARCHITECTURE.md index 5134794c036..d4df7399b77 100644 --- a/packages/spec/docs/SYNC_ARCHITECTURE.md +++ b/packages/spec/docs/SYNC_ARCHITECTURE.md @@ -95,9 +95,9 @@ ten-stage pipeline, get no error, and get no execution. **What to use instead — layer by layer, and one honest gap:** - **Scheduled pull from an external system** — the target-side binding - `mapping.connectorSource` with a `job` for the cadence (pulled when a job drives it; - the job stage that schedules it has not landed — see - [Data sync is defined on the target](#data-sync-is-defined-on-the-target)). + `mapping.connectorSource`, pulled by a `job` whose `pull: { mapping }` names it, on + the job's schedule — see + [Data sync is defined on the target](#data-sync-is-defined-on-the-target). - **Per-field value conversion on import** — `mapping.fieldMapping[].transform` (`data/mapping.zod.ts`): a string enum (`none` / `constant` / `map` / `split` / `join` / `lookup`) with its settings in `params`, applied row by row by the REST @@ -229,9 +229,19 @@ Complete, production-grade integration with external systems. Includes authentic > incremental pull over a newest-first paged endpoint moves its starting point past > the pages it never read. Point `connectorSource` at an endpoint that answers the > whole (incremental) set in one response. -> - ⚠️ **Nothing schedules a pull yet.** A `job` drives it, and that stage has not -> landed, so the binding alone moves no rows — `connectorSource`'s own description -> says so. It is not a lint warning: every key of the binding is `live`. +> - **A `job` drives the pull.** The binding alone moves no rows: a `job` whose +> `pull: { mapping }` names the mapping pulls it on the job's `schedule` — no code, a +> run form of its own, refused beside `body` or `handler`. `defineStack` (and so +> `os validate`) refuses a `pull` naming a mapping the stack does not declare, or one +> with no `connectorSource`. Each run calls the `automation` service's +> `pullConnectorSource`; a refused pull records the run `failed` (retried per the +> job's `retryPolicy`), a pull whose rows the import runner refused records it +> `degraded` with the counts. +> - **The organization it writes as.** The pull runs as the job's declared +> `organization` (`{ isSystem: true, tenantId }`), judged at bind by the posture rule +> scheduled flows use: required under the `isolated` tenancy posture (a job with none +> is not scheduled), optional under `group` (undeclared, a tenant-scoped write is +> refused), not required under `single`. > Already authored the retired keys? `os migrate meta --from 17` lists the mechanical edits. ### Use Cases @@ -373,7 +383,7 @@ mostly answers "which surface", and — for the two questions that used to route | Do you need real-time webhooks? | **Outbound:** the stack's top-level `webhooks:` collection (`src/automation/webhook.zod.ts`) — **not** L3: a connector's nested `webhooks` was never delivered and is retired (ADR-0049) | | Do you need advanced authentication (OAuth2, SAML)? | **Yes** → L3 (Connector) | | Do you need retry policies and circuit breaking? | **Retry: yes, L3.** `retryConfig` is executed at the platform's one outbound call (ADR-0049 ruled `实现`) — backoff shape, attempt count, retryable statuses, network-error retry and a per-attempt `requestTimeoutMs`. **Circuit breaking: no level provides it** — `health.circuitBreaker` was retired (ADR-0049) because no breaker ever opened; implement it in the connector provider or an upstream gateway. Outbound **rate limiting** is not a reason to pick any level either: no level provides it (#4911); throttle at the provider or gateway | -| Is it a simple pull from an external system into a local object? | The target-side binding: a `mapping` with `connectorSource` over a `rest` / `openapi` connector, cadence from a `job` — **pulled when a job drives it; the job stage has not landed** ([above](#data-sync-is-defined-on-the-target)) | +| Is it a simple pull from an external system into a local object? | The target-side binding: a `mapping` with `connectorSource` over a `rest` / `openapi` connector, pulled on a cadence by a `job` whose `pull: { mapping }` names it ([above](#data-sync-is-defined-on-the-target)) | | Are you building a data warehouse pipeline? | The extraction half is that same pull binding; the warehouse-side transformation is the warehouse's own tooling. There is no ObjectStack pipeline protocol (#6414) | | Are you integrating with an enterprise system? | **Yes** → L3 (Connector) | | Do you need client-side offline sync? | Not this layering — and note `ui/offline.zod.ts` was itself retired at #4988 for having no carrier key | @@ -394,8 +404,8 @@ instance with simple `auth` whose actions flows call. External API → L3 Connector → ObjectStack → (warehouse's own ELT) ``` The second arrow used to read `ObjectStack → L2 ETL → Data Warehouse`, and that hop -never executed. Land the data through a connector (the target-side pull binding, once -a job drives it), then transform it with a tool that actually runs — the +never executed. Land the data through a connector (the target-side pull binding, driven +by a job's `pull`), then transform it with a tool that actually runs — the warehouse's own ELT, a `flow`, or a scheduled job. --- @@ -425,11 +435,12 @@ const pipeline: ETLPipeline = { ``` **After** — split it by which half has a runtime. The extraction half is the -target-side pull binding (pulled when a job drives it; the job stage has not landed): +target-side pull binding, pulled on a cadence by a job whose `pull` names it: ```typescript import type { Connector } from '@objectstack/spec/integration'; import type { Mapping } from '@objectstack/spec/data'; +import type { Job } from '@objectstack/spec/system'; const orders: Connector = { name: 'orders', @@ -453,6 +464,13 @@ const ordersPull: Mapping = { watermark: { field: 'updated_at', param: 'updated_since' }, }, }; + +// The cadence: a job whose `pull` names the mapping — no code. +const ordersPullHourly: Job = { + name: 'orders_pull_hourly', + schedule: { type: 'cron', expression: '0 * * * *', timezone: 'UTC' }, + pull: { mapping: 'orders_pull' }, +}; ``` The `transformations` half has no runtime, and never did. Aggregations, joins and diff --git a/packages/spec/dropped-refinements.baseline.json b/packages/spec/dropped-refinements.baseline.json index dba54ed198f..6b4e590c64a 100644 --- a/packages/spec/dropped-refinements.baseline.json +++ b/packages/spec/dropped-refinements.baseline.json @@ -3,7 +3,7 @@ "measured": { "zod": "4.4.3", "publishedSchemasWithDroppedRefinements": 218, - "droppedRefinementSites": 660, + "droppedRefinementSites": 665, "refinementSitesThatDidProject": 369, "refinementSitesWithNoJsonFormToCompare": 0 }, @@ -82,6 +82,7 @@ "manifest.flows.element.nodes.element.in.waitEventConfig", "manifest.hooks.element", "manifest.hooks.element.object", + "manifest.jobs.element", "manifest.jobs.element.schedule.options[0].timezone", "manifest.navigationContributions.element.items.element.lazy.options[0]", "manifest.objectExtensions.element", @@ -191,6 +192,7 @@ "data.options[1].manifest.flows.element.nodes.element.in.waitEventConfig", "data.options[1].manifest.hooks.element", "data.options[1].manifest.hooks.element.object", + "data.options[1].manifest.jobs.element", "data.options[1].manifest.jobs.element.schedule.options[0].timezone", "data.options[1].manifest.objectExtensions.element", "data.options[1].manifest.objects.element", @@ -285,6 +287,7 @@ "options[1].manifest.flows.element.nodes.element.in.waitEventConfig", "options[1].manifest.hooks.element", "options[1].manifest.hooks.element.object", + "options[1].manifest.jobs.element", "options[1].manifest.jobs.element.schedule.options[0].timezone", "options[1].manifest.objectExtensions.element", "options[1].manifest.objects.element", @@ -339,6 +342,7 @@ "data.packages.element.options[1].manifest.flows.element.nodes.element.in.waitEventConfig", "data.packages.element.options[1].manifest.hooks.element", "data.packages.element.options[1].manifest.hooks.element.object", + "data.packages.element.options[1].manifest.jobs.element", "data.packages.element.options[1].manifest.jobs.element.schedule.options[0].timezone", "data.packages.element.options[1].manifest.objectExtensions.element", "data.packages.element.options[1].manifest.objects.element", @@ -1107,6 +1111,7 @@ }, "system/Job": { "sites": [ + "", "schedule.options[0].timezone" ] }, diff --git a/packages/spec/export-origins/contracts.json b/packages/spec/export-origins/contracts.json index 425c26c7607..e8c01bdd981 100644 --- a/packages/spec/export-origins/contracts.json +++ b/packages/spec/export-origins/contracts.json @@ -61,6 +61,9 @@ "CheckNamespaceInput": "src/contracts/package-service.ts#CheckNamespaceInput (interface)", "CheckNamespaceResult": "src/contracts/package-service.ts#CheckNamespaceResult (interface)", "ClusterCallContext": "src/contracts/cluster-service.ts#ClusterCallContext (interface)", + "ConnectorSourcePullRequest": "src/contracts/automation-service.ts#ConnectorSourcePullRequest (interface)", + "ConnectorSourcePullResult": "src/contracts/automation-service.ts#ConnectorSourcePullResult (interface)", + "ConnectorSourcePullSummary": "src/contracts/automation-service.ts#ConnectorSourcePullSummary (interface)", "CoreServiceContract": "src/contracts/core-service-contracts.ts#CoreServiceContract (type)", "CoreServiceContracts": "src/contracts/core-service-contracts.ts#CoreServiceContracts (interface)", "CounterIncrOptions": "src/contracts/cluster-service.ts#CounterIncrOptions (interface)", diff --git a/packages/spec/liveness/job.json b/packages/spec/liveness/job.json index 4a44494d5f4..7cb0c6cf550 100644 --- a/packages/spec/liveness/job.json +++ b/packages/spec/liveness/job.json @@ -1,6 +1,6 @@ { "type": "job", - "_note": "JobSchema. The file-authored path is healthy: `defineStack({ jobs })` → `AppPlugin.start`'s `kernel:ready` hook → IJobService.schedule → the service-job adapters honor every schedule shape (`CronJobAdapter.schedule` / `DbJobAdapter.schedule`) and `runWithPolicy` enforces retryPolicy/timeout (#3494 — these used to be parsed-but-ignored). `retryPolicy` here is the ENFORCED spelling ({maxRetries, backoffMs, backoffMultiplier}); do not confuse it with the datasource `retryPolicy`, which is dead and spells its delay differently. TYPE-LEVEL GAP CLOSED 2026-08-02 (#4509) by CLOSING THE DOOR, not building a bridge: `job` was registered `allowRuntimeCreate: true` while only the compiled bundle's `jobs` ever reached the scheduler, so a Studio-created job saved cleanly and never ran. Unlike the webhook (#3461) and email_template (#4509 item 1) disconnects, this one could not be bridged: `handler` names a function in the compiled bundle's function table (`collectBundleFunctions`), which a runtime writer does not have and cannot name — the missing piece is a handler-binding design, not an ingestion path. So `allowRuntimeCreate` AND `allowOrgOverride` are now both false (metadata-plugin.zod.ts, with the rationale block), leaving `*.job.ts` / `defineStack({ jobs })` as the supported doors. The kind stays registered: its file loader is genuinely consumed, so it still passes the ADR-0088 admission test. Seeded 2026-08-01. 2026-08-28 (commit 8cb96ec41): every LOCAL citation re-anchored to its consuming symbol — and this file is its own best argument for doing so. The 2026-08-02 note recorded that the SEEDED lines had already drifted ~25 lines and were restamped with fresh numbers; twenty-six days later every one of those fresh numbers had drifted again, this time by ~70-100 lines, onto `} else {`, `try {`, a bare `}` and a line about registering ACTIONS. Restamping is not a fix for line rot, it is the same claim with a newer date — which is the case this whole worklist rests on. The two objectui-cited entries (`label`, `description`) are left BYTE-FOR-BYTE UNTOUCHED: their evidence and producer are pinned at `objectui @aeb8424b`, a commit this container cannot reproduce, and foreign anchors are never collected by the scanner. 2026-10-03 (#21489): a job's sandboxed `body` is scheduled by the binder's job half (`scheduleAppArtifactJobs`, `packages/runtime/src/app-artifact-handlers.ts`) on every door that brings an artifact in — the boot and install-local (install and rehydrate) — and the install-local door refuses a package whose enabled job has no `body`, the one shape no JSON door can run.", + "_note": "JobSchema. The file-authored path is healthy: `defineStack({ jobs })` → `AppPlugin.start`'s `kernel:ready` hook → IJobService.schedule → the service-job adapters honor every schedule shape (`CronJobAdapter.schedule` / `DbJobAdapter.schedule`) and `runWithPolicy` enforces retryPolicy/timeout (#3494 — these used to be parsed-but-ignored). `retryPolicy` here is the ENFORCED spelling ({maxRetries, backoffMs, backoffMultiplier}); do not confuse it with the datasource `retryPolicy`, which is dead and spells its delay differently. TYPE-LEVEL GAP CLOSED 2026-08-02 (#4509) by CLOSING THE DOOR, not building a bridge: `job` was registered `allowRuntimeCreate: true` while only the compiled bundle's `jobs` ever reached the scheduler, so a Studio-created job saved cleanly and never ran. Unlike the webhook (#3461) and email_template (#4509 item 1) disconnects, this one could not be bridged: `handler` names a function in the compiled bundle's function table (`collectBundleFunctions`), which a runtime writer does not have and cannot name — the missing piece is a handler-binding design, not an ingestion path. So `allowRuntimeCreate` AND `allowOrgOverride` are now both false (metadata-plugin.zod.ts, with the rationale block), leaving `*.job.ts` / `defineStack({ jobs })` as the supported doors. The kind stays registered: its file loader is genuinely consumed, so it still passes the ADR-0088 admission test. Seeded 2026-08-01. 2026-08-28 (commit 8cb96ec41): every LOCAL citation re-anchored to its consuming symbol — and this file is its own best argument for doing so. The 2026-08-02 note recorded that the SEEDED lines had already drifted ~25 lines and were restamped with fresh numbers; twenty-six days later every one of those fresh numbers had drifted again, this time by ~70-100 lines, onto `} else {`, `try {`, a bare `}` and a line about registering ACTIONS. Restamping is not a fix for line rot, it is the same claim with a newer date — which is the case this whole worklist rests on. The two objectui-cited entries (`label`, `description`) are left BYTE-FOR-BYTE UNTOUCHED: their evidence and producer are pinned at `objectui @aeb8424b`, a commit this container cannot reproduce, and foreign anchors are never collected by the scanner. 2026-10-03 (#21489): a job's sandboxed `body` is scheduled by the binder's job half (`scheduleAppArtifactJobs`, `packages/runtime/src/app-artifact-handlers.ts`) on every door that brings an artifact in — the boot and install-local (install and rehydrate) — and the install-local door refuses a package whose enabled job has no `body`, the one shape no JSON door can run. 2026-10-04 (#20281 stage ③, rulings Q1-B + Q2-O1): two keys land — `pull`, the declarative run form (`{ mapping }`: the platform pulls that mapping's `connectorSource` on the job's schedule through the `automation` service's contract method `pullConnectorSource`, bound by the same binder's job half), and `organization`, the organization every run form runs as, judged at bind by the scheduled flows' posture rule (`resolveScheduledWorkPolicy`).", "props": { "name": { "status": "live", @@ -76,6 +76,27 @@ }, "note": "Drilled on the day the key landed: `body` is `ScriptBodySchema` behind a slot refinement, and its five keys do not share one verdict — `timeoutMs` is refused on a job (at parse, and at bind for a body that never parsed) while the other four are live since #21489's binder." }, + "pull": { + "children": { + "mapping": { + "status": "live", + "verifiedAt": "2026-10-04", + "evidenceScope": "in-repo", + "evidence": "packages/runtime/src/app-artifact-handlers.ts#judgeJobPull (the name is parsed against `JobSchema.pull` and resolved against the artifact's mappings; a missing mapping, or one with no `connectorSource`, is refused at bind and the job is not scheduled); packages/services/service-automation/src/engine.ts#AutomationEngine (`pullConnectorSource` — the `IAutomationService` contract method every run of the job calls with the name); packages/services/service-automation/src/connector-pull.ts#pullConnectorSource (reads the named mapping through the protocol and pulls its `connectorSource` through the import runner); packages/spec/src/stack.zod.ts#collectJobPullMappingErrors (`defineStack`, and so `os validate`, refuses a name the stack's mappings do not resolve, or one whose mapping declares no `connectorSource`)", + "producer": "packages/runtime/src/app-artifact-handlers.ts#scheduleAppArtifactJobs (each run hands the judged name to `service.pullConnectorSource({ mapping, context })`, the `automation` service resolved through the registry); packages/services/service-automation/src/plugin.ts#AutomationServicePlugin (`init()` attaches the executor to the engine with `setConnectorPullSource` before registering the engine as the `automation` service)", + "note": "LIVE on the day it landed (#20281 stage ③, ruling Q1-B): the name a job pulls. The binder schedules a `pull` job on every door that brings an artifact in (the boot, install-local install and rehydrate) and each run calls the contract method; a refused pull rejects the run (`failed`, retried per `retryPolicy`), a pull whose rows the import runner refused resolves `degraded`." + } + }, + "note": "The declarative run form (ruling Q1-B), a strict one-key object drilled on the day it landed. Exclusive with `body` and `handler` (refused at parse; that refinement is a recorded dropped site in `dropped-refinements.baseline.json`, so the published JSON Schema is wider than the parse here)." + }, + "organization": { + "status": "live", + "verifiedAt": "2026-10-04", + "evidenceScope": "in-repo", + "evidence": "packages/runtime/src/app-artifact-handlers.ts#scheduleAppArtifactJobs (`resolveJobOrganization(job)` — under `policy.requiresActingOrganization` a job declaring none is NOT scheduled, at `error`; under `runOwnership: 'per-record'` an undeclared one is named once at `warn`; every run of every form carries `jobExecutionContext(organization)`); packages/runtime/src/app-artifact-handlers.ts#jobExecutionContext (`{ isSystem: true, tenantId }`, handed to a `pull` as `context` and to a `handler` as `executionContext`); packages/runtime/src/sandbox/body-runner.ts#buildJobSandboxContext (a `body`'s `ctx.api` envelope carries it as `tenantId`); packages/objectql/src/engine.ts#resolveSystemInsertOrganization (`carriesOrganization(execCtx?.tenantId)` — the tenancy guard a declared organization satisfies, instead of the `walled-posture` refusal)", + "producer": "packages/types/src/env.ts#resolveScheduledWorkPolicy (the deployment posture the binder judges the declaration against — `requiresActingOrganization` / `runOwnership`, read from the environment at bind; the same resolver scheduled flows bind by)", + "note": "LIVE on the day it landed (#20281 stage ③, ruling Q2-O1: the maintainer's 2026-09-08 ruling on scheduled work across organizations, extended from scheduled flows to jobs). The value shape is `ScheduleOrganizationSchema` (`automation/schedule-organization.zod.ts`), reused by reference. Judged at BIND, never at authoring — whether it is required depends on the posture and the scheduled-work switch, which no authoring door can read." + }, "retryPolicy": { "status": "live", "verifiedAt": "2026-10-03", diff --git a/packages/spec/liveness/mapping.json b/packages/spec/liveness/mapping.json index 49cc72e394b..ee152cfef4a 100644 --- a/packages/spec/liveness/mapping.json +++ b/packages/spec/liveness/mapping.json @@ -48,7 +48,7 @@ "status": "live", "verifiedAt": "2026-10-01", "evidence": "packages/services/service-automation/src/connector-pull.ts#pullConnectorSource (reads the binding through the protocol's `getMetaItem({ type: 'mapping', name })`, makes ONE action call on the declared connector, and writes through `@objectstack/core`'s `runImport` — the import door's runner)", - "note": "Seeded 2026-09-30 `planned`, contract-first — the target-side half of the connector-attached sync family, ruled ENFORCE on the maintainer's criterion with the definition MOVED here from the retired `connector.syncConfig` / `connector.fieldMappings` (both `dead` tombstone rows in connector.json): every mainstream platform binds a sync to its TARGET, and this type already carries the executed target half (`targetObject`, `fieldMapping`, `mode`, `upsertKey`, all live). `connectorSource` adds only where the rows come from — a `rest` / `openapi` connector instance, the action that reads them and an optional timestamp watermark; version 1 is a one-way pull, full or incremental, credentials are the instance's own ADR-0097 static `auth`, and the cadence is a `job` (no schedule key, by the 2026-09-10 ruling). LIVE since 2026-10-01 (#20919, stage ②): `@objectstack/service-automation`'s connector sync executor reads every key below and writes through the import runner, which moved to `@objectstack/core` so the executor writes through the same runner, coercion and row verdicts as the import door. `authorWarn` was kept at that re-grade and dropped on 2026-10-01 (#21127): a `live` row carries no author warning — the author-side lint has no verdict for one and throws its ledger-integrity error, so `os validate` and `os lint` exited 1 on every stack that authored the binding (`check:liveness` refuses the combination now). The caveat the warning carried — nothing schedules a pull until the `job` stage lands (stage ③), so a pull runs only when a job drives it — is on the key's `.describe()`, where an author reads it; stage ③ removes it there.", + "note": "Seeded 2026-09-30 `planned`, contract-first — the target-side half of the connector-attached sync family, ruled ENFORCE on the maintainer's criterion with the definition MOVED here from the retired `connector.syncConfig` / `connector.fieldMappings` (both `dead` tombstone rows in connector.json): every mainstream platform binds a sync to its TARGET, and this type already carries the executed target half (`targetObject`, `fieldMapping`, `mode`, `upsertKey`, all live). `connectorSource` adds only where the rows come from — a `rest` / `openapi` connector instance, the action that reads them and an optional timestamp watermark; version 1 is a one-way pull, full or incremental, credentials are the instance's own ADR-0097 static `auth`, and the cadence is a `job` (no schedule key, by the 2026-09-10 ruling). LIVE since 2026-10-01 (#20919, stage ②): `@objectstack/service-automation`'s connector sync executor reads every key below and writes through the import runner, which moved to `@objectstack/core` so the executor writes through the same runner, coercion and row verdicts as the import door. `authorWarn` was kept at that re-grade and dropped on 2026-10-01 (#21127): a `live` row carries no author warning — the author-side lint has no verdict for one and throws its ledger-integrity error, so `os validate` and `os lint` exited 1 on every stack that authored the binding (`check:liveness` refuses the combination now). The caveat the warning carried — nothing scheduled a pull until the `job` stage landed (stage ③), so a pull ran only when a job drove it — was on the key's `.describe()`, where an author reads it. 2026-10-04 (#20281 stage ③): a `job` whose `pull` names this mapping drives it (`JobSchema.pull`, through the `automation` service's `pullConnectorSource`), and the describe now says that instead.", "children": { "connector": { "status": "live", diff --git a/packages/spec/liveness/state-counts/job.md b/packages/spec/liveness/state-counts/job.md index 80dfd4bbe43..c60b6b9a0ed 100644 --- a/packages/spec/liveness/state-counts/job.md +++ b/packages/spec/liveness/state-counts/job.md @@ -12,4 +12,4 @@ committed anywhere: `check:liveness` sums the shards when it reads them. | Type | live | exp | elsewhere | dead | planned | classified | |---|---|---|---|---|---|---| -| `job` | 19 | 0 | 0 | 1 | 1 | 21 | +| `job` | 21 | 0 | 0 | 1 | 1 | 23 | diff --git a/packages/spec/src/contracts/automation-service.ts b/packages/spec/src/contracts/automation-service.ts index b9c6f6042c5..90bb4566012 100644 --- a/packages/spec/src/contracts/automation-service.ts +++ b/packages/spec/src/contracts/automation-service.ts @@ -18,6 +18,7 @@ import type { ExecutionLog, ExecutionStatus, FlowRunSummary } from '../automatio import type { ActionDescriptor } from '../automation/node-executor.zod'; import type { ConnectorDescriptor } from '../integration/connector-descriptor'; import type { ConversionNotice, ConversionConflictNotice } from '../conversions/types'; +import type { ExecutionContext } from '../kernel/execution-context.zod'; /** * Context passed to a flow/script execution @@ -643,6 +644,65 @@ export interface RunListResult { hasMore: boolean; } +/** + * What {@link IAutomationService.pullConnectorSource} is asked to pull: one + * `mapping` by name, under the caller's execution context. + * + * The caller is a job (`JobSchema.pull`): the job binder builds the context + * from the job's declared `organization` — `{ isSystem: true, tenantId }`, or + * `{ isSystem: true }` where it declares none — so the pull's target read and + * its writes carry the organization the job runs as (ruling Q2-O1 on the + * connector-sync card). + */ +export interface ConnectorSourcePullRequest { + /** Name of the `mapping` whose `connectorSource` is pulled. */ + mapping: string; + /** The execution context the target read and the writes run under. */ + context?: ExecutionContext; + /** The environment the mapping and the target are read in, when the caller is environment-scoped. */ + environmentId?: string; +} + +/** + * The import runner's tallies for the rows one pull wrote — the counters of + * the runner's own run summary, which carries more (per-row results); this is + * the part a caller of the contract may read. + */ +export interface ConnectorSourcePullSummary { + /** Rows handed to the runner. */ + total: number; + /** Rows the runner reached. */ + processed: number; + created: number; + updated: number; + skipped: number; + /** Rows the runner REFUSED (a row verdict, not a refused pull). */ + errors: number; + /** Rows written (`created + updated`). */ + ok: number; + cancelled: boolean; +} + +/** + * What one {@link IAutomationService.pullConnectorSource} call did. A refused + * pull never answers this — it rejects before anything is written. + */ +export interface ConnectorSourcePullResult { + mapping: string; + targetObject: string; + connector: string; + action: string; + /** + * The incremental half, when the mapping's `connectorSource.watermark` is + * declared: `from` is the starting point sent as `query[param]` (the + * highest value already stored in `target`), `undefined` on a first pull. + */ + watermark?: { field: string; target: string; param: string; from: unknown }; + /** Records the one action response carried. */ + pulled: number; + summary: ConnectorSourcePullSummary; +} + export interface IAutomationService { /** * Execute a named flow or script @@ -840,6 +900,31 @@ export interface IAutomationService { */ getConnectorDescriptors?(): ConnectorDescriptor[]; + /** + * Pull one `mapping`'s `connectorSource` and write the records through the + * import runner — the connector sync executor, on the automation service. + * + * The door a job's `pull` run form binds through (`JobSchema.pull`, ruling + * Q1-B on the connector-sync card): the job binder resolves the + * `automation` service and calls this on every run, with the context the + * job's `organization` builds. Before it existed the executor was reachable + * only as an instance method of the automation PLUGIN, which no other + * package holds. + * + * Rejects — before anything is written — when the binding cannot be + * honoured (no such mapping, no `connectorSource`, a connector that is not + * a declared `rest` / `openapi` instance, an upstream `ok: false`, …), with + * an error carrying an ADR-0112 `code` + `status` and a `reason` + * discriminator. A pull whose rows the import runner refused RESOLVES, and + * says so in `summary.errors`. + * + * Optional for the reason {@link getConnectorDescriptors} is: a connector + * registry is a capability of the flow-engine implementation, not of every + * automation slot. A caller that finds it absent reports that nothing can + * pull, rather than calling. + */ + pullConnectorSource?(request: ConnectorSourcePullRequest): Promise; + /** * Per-flow deployment + binding state, for operator surfaces. * diff --git a/packages/spec/src/data/mapping-connector-source.test.ts b/packages/spec/src/data/mapping-connector-source.test.ts index 85851561944..129b573b90c 100644 --- a/packages/spec/src/data/mapping-connector-source.test.ts +++ b/packages/spec/src/data/mapping-connector-source.test.ts @@ -18,8 +18,9 @@ * the binding, so the ledger rows are `live` and carry no `authorWarn` * (`liveness/mapping.json`) — a warned `live` row made the author-side lint * throw, so `os validate` / `os lint` crashed on every stack authoring the - * binding (#21127). Nothing schedules a pull until the `job` stage lands, and - * the key's own description says so. The last block pins both. + * binding (#21127). A pull runs when a `job` whose `pull` names the mapping + * drives it (`JobSchema.pull`), and the key's own description says so. The + * last block pins both. */ import fs from 'node:fs'; @@ -177,7 +178,7 @@ describe('mapping.connectorSource — the closed door and what it does not carry }); }); -describe('mapping.connectorSource — executed when a job drives it, scheduled by nothing yet: the ledger and the schema say so', () => { +describe('mapping.connectorSource — executed when a job\'s `pull` drives it: the ledger and the schema say so', () => { const LEDGER = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../../liveness/mapping.json'); it('every key of the binding is `live`, citing the executor, and no row of it opts into an author warning', () => { @@ -218,7 +219,10 @@ describe('mapping.connectorSource — executed when a job drives it, scheduled b it('the schema says it too, where an author reads it — including where the watermark is read from and the one-response limit', () => { const shape = (MappingSchema as unknown as { shape: Record unknown }> }).shape; const source = shape.connectorSource!; - expect(source.description).toContain('Pulled when a job drives it; nothing schedules it yet'); + expect(source.description).toContain('Pulled when a job drives it'); + expect(source.description).toContain('`pull: { mapping }`'); + expect(source.description).toContain('the binding alone moves no rows'); + expect(source.description).not.toContain('nothing schedules'); const inner = (source.unwrap!() as { shape: Record unknown }> }).shape; const watermark = inner.watermark!; expect(watermark.description).toContain('ONE action call and reads ONE response'); diff --git a/packages/spec/src/data/mapping.zod.ts b/packages/spec/src/data/mapping.zod.ts index 0a9183b80cf..0cc71aefaae 100644 --- a/packages/spec/src/data/mapping.zod.ts +++ b/packages/spec/src/data/mapping.zod.ts @@ -363,17 +363,16 @@ export const MappingSchema = lazySchema(() => strictObject({ * - **no delete or conflict policy** — a pull writes through this mapping's * `mode` / `upsertKey`, and nothing else is claimed. * - * EXECUTED WHEN A JOB DRIVES IT: `@objectstack/service-automation`'s - * connector sync executor (`pullConnectorSource`) reads every key here and - * writes through the import runner, so the liveness ledger grades every key - * `live`. Nothing schedules a pull yet — the `job` that drives it is the - * next stage — and that is the one caveat an author must read before - * writing the binding, so the `.describe()` below carries it. It is not a - * ledger warning: a `live` row carries no `authorWarn` (the author-side lint - * has no verdict for one, and `check:liveness` refuses it). A pull makes ONE - * action call and reads ONE response (see `watermark`). `sourceFormat` keeps - * governing the manual import door; a pulled row is the connector's JSON - * record. + * EXECUTED WHEN A JOB PULLS IT: `@objectstack/service-automation`'s + * connector sync executor (`pullConnectorSource`, on the `automation` + * service's contract) reads every key here and writes through the import + * runner, so the liveness ledger grades every key `live`. A `job` whose + * `pull` names this mapping (`JobSchema.pull`) drives it on the job's + * schedule, as the organization the job declares; the binding alone moves + * no rows, which the `.describe()` below says where an author reads it. A + * pull makes ONE action call and reads ONE response (see `watermark`). + * `sourceFormat` keeps governing the manual import door; a pulled row is + * the connector's JSON record. */ connectorSource: strictObject({ surface: 'this mapping’s connector source', @@ -427,8 +426,8 @@ export const MappingSchema = lazySchema(() => strictObject({ ), }).optional().describe( 'Pull binding: the rest/openapi connector this mapping pulls rows from (one-way, full or ' - + 'timestamp-incremental; a `job` sets the cadence). Pulled when a job drives it; nothing schedules it yet, ' - + 'so the binding alone moves no rows — schedule the pull with a `job` once a job can drive one', + + 'timestamp-incremental; a `job` sets the cadence). Pulled when a job drives it — a `job` whose ' + + '`pull: { mapping }` names this mapping, on the job\'s schedule; the binding alone moves no rows', ), // `extractQuery`, `errorPolicy` and `batchSize` were removed in 17.0.0 diff --git a/packages/spec/src/integration/connector-sync-retirement.test.ts b/packages/spec/src/integration/connector-sync-retirement.test.ts index 9dd7a825bd4..9a13d1f46b1 100644 --- a/packages/spec/src/integration/connector-sync-retirement.test.ts +++ b/packages/spec/src/integration/connector-sync-retirement.test.ts @@ -117,15 +117,17 @@ describe('connector sync retirement — the tombstones', () => { }); } - it('the `syncConfig` prescription is honest about the binding it points at: pulled when a job drives it, scheduled by nothing yet', () => { + it('the `syncConfig` prescription is honest about the binding it points at: pulled when a job\'s `pull` names the mapping, never on its own', () => { // The pull executor reads the target-side binding (its ledger rows are - // `live`), but nothing schedules a pull until the `job` stage lands (the - // binding's own description says so), so a prescription that sent the author - // there as if it ran on its own would be the defect this retirement - // removes, moved one type over. + // `live`), and a pull runs only when a `job` whose `pull` names the mapping + // drives it (`JobSchema.pull`), so a prescription that sent the author + // there as if the binding ran on its own would be the defect this + // retirement removes, moved one type over. const message = issueAt(ConnectorSchema.safeParse({ ...WELL_FORMED, syncConfig: {} }), 'syncConfig')!.message; expect(message).toContain('Its pull runs when a `job` drives it'); - expect(message).toContain('nothing schedules one yet'); + expect(message).toContain('`pull: { mapping }`'); + expect(message).toContain('the binding alone moves no rows'); + expect(message).not.toContain('nothing schedules'); }); it('refuses EVERY value — an empty block, an empty list, null and a scalar included', () => { diff --git a/packages/spec/src/integration/connector.zod.ts b/packages/spec/src/integration/connector.zod.ts index b9e89f48269..d1902dadc01 100644 --- a/packages/spec/src/integration/connector.zod.ts +++ b/packages/spec/src/integration/connector.zod.ts @@ -228,8 +228,8 @@ const SYNC_CONFIG_RETIRED = + 'Delete the key; the `DataSyncConfig` shape leaves with it. A sync is defined on its TARGET: ' + 'a `mapping` (`targetObject`, `fieldMapping`, `mode`, `upsertKey`) whose `connectorSource` ' + 'names the `rest` or `openapi` connector it pulls from, the read action and an optional ' - + 'timestamp `watermark`, with a `job` for the cadence. Its pull runs when a `job` drives it; ' - + 'nothing schedules one yet, so the binding alone moves no rows. ' + + 'timestamp `watermark`, with a `job` for the cadence. Its pull runs when a `job` drives it — ' + + 'a `job` whose `pull: { mapping }` names that mapping; the binding alone moves no rows. ' + 'Run `os migrate meta --from 17` to list the mechanical edits for existing sources; apply them by hand.'; /** diff --git a/packages/spec/src/migrations/entries/semantic/18.connector-sync-keys-retired.ts b/packages/spec/src/migrations/entries/semantic/18.connector-sync-keys-retired.ts index 0aca330a413..2e87ce8ea9f 100644 --- a/packages/spec/src/migrations/entries/semantic/18.connector-sync-keys-retired.ts +++ b/packages/spec/src/migrations/entries/semantic/18.connector-sync-keys-retired.ts @@ -24,8 +24,8 @@ export const entry: SemanticMigration = { + 'instance it pulls from (`connector`), the action that reads the records (`action`, with a ' + 'fixed `input` and a `recordsPath`) and, for a timestamp-incremental pull, a `watermark` ' + '(`field` on the record, `param` on the request); a `job` sets the cadence. The pull executor ' - + 'reads the binding when a `job` drives it; nothing schedules a pull yet, so the binding alone ' - + 'moves no rows.', + + 'reads the binding when a `job` drives it — a `job` whose `pull: { mapping }` names the mapping, ' + + 'on the job\'s schedule; the binding alone moves no rows.', reason: 'The D2 conversion `connector-sync-keys-removed` deletes `syncConfig` and ' + '`fieldMappings` from every connector, stack entry and stored connector row, one notice per ' + 'key, and the delete is lossless: no engine ever ran a connector-attached sync or moved a ' @@ -48,6 +48,6 @@ export const entry: SemanticMigration = { + 'DataSyncConfig, SyncStrategy, ConnectorConflictResolution or ConnectorFieldMapping or ' + 'their schemas. Every connector registers and dispatches its actions exactly as it did ' + 'before the upgrade. Each sync the author still wants is a `mapping` whose ' - + '`connectorSource` names a `rest` or `openapi` connector instance, with a `job` chosen for ' - + 'its cadence.', + + '`connectorSource` names a `rest` or `openapi` connector instance, with a `job` whose `pull` ' + + 'names that mapping chosen for its cadence.', }; diff --git a/packages/spec/src/migrations/registry.ts b/packages/spec/src/migrations/registry.ts index 69dffdc5eb9..ce2b787d3fa 100644 --- a/packages/spec/src/migrations/registry.ts +++ b/packages/spec/src/migrations/registry.ts @@ -8793,8 +8793,8 @@ const step18: MigrationStep = { + 'instance it pulls from (`connector`), the action that reads the records (`action`, with a ' + 'fixed `input` and a `recordsPath`) and, for a timestamp-incremental pull, a `watermark` ' + '(`field` on the record, `param` on the request); a `job` sets the cadence. The pull executor ' - + 'reads the binding when a `job` drives it; nothing schedules a pull yet, so the binding alone ' - + 'moves no rows.', + + 'reads the binding when a `job` drives it — a `job` whose `pull: { mapping }` names the mapping, ' + + 'on the job\'s schedule; the binding alone moves no rows.', reason: 'The D2 conversion `connector-sync-keys-removed` deletes `syncConfig` and ' + '`fieldMappings` from every connector, stack entry and stored connector row, one notice per ' + 'key, and the delete is lossless: no engine ever ran a connector-attached sync or moved a ' @@ -8817,8 +8817,8 @@ const step18: MigrationStep = { + 'DataSyncConfig, SyncStrategy, ConnectorConflictResolution or ConnectorFieldMapping or ' + 'their schemas. Every connector registers and dispatches its actions exactly as it did ' + 'before the upgrade. Each sync the author still wants is a `mapping` whose ' - + '`connectorSource` names a `rest` or `openapi` connector instance, with a `job` chosen for ' - + 'its cadence.', + + '`connectorSource` names a `rest` or `openapi` connector instance, with a `job` whose `pull` ' + + 'names that mapping chosen for its cadence.', }, // ADR-0049 enforce-or-remove — the D3 entry of the connector triggers family: // `connector.triggers`, the whole `ConnectorTrigger` array, retired as one diff --git a/packages/spec/src/stack.zod.ts b/packages/spec/src/stack.zod.ts index 0029bf064db..df494457abc 100644 --- a/packages/spec/src/stack.zod.ts +++ b/packages/spec/src/stack.zod.ts @@ -2825,6 +2825,51 @@ function collectHydratedInlineColumnErrors(config: ObjectStackDefinition): strin return errors; } +/** + * A job's `pull.mapping` must name a mapping this stack declares, and that + * mapping must carry a `connectorSource` (ruling Q1-B on the connector-sync + * card: "`os validate` checks the mapping name"). + * + * Both halves are refused here rather than at the first tick, where the + * automation service would reject the run (`mapping_not_found` / + * `no_connector_source`): a job is package-authored, and a reference its own + * stack cannot resolve is an authoring defect the package would otherwise ship. + * Resolved against the stack's own `mappings`, like every mapping reference in + * this function. A disabled job is judged too — `enabled: false` turns a job + * off, it does not make a dangling name correct. + */ +function collectJobPullMappingErrors(config: ObjectStackDefinition): string[] { + const errors: string[] = []; + const jobs = Array.isArray(config.jobs) ? config.jobs : []; + if (jobs.length === 0) return errors; + const mappings = new Map(); + for (const m of Array.isArray(config.mappings) ? config.mappings : []) { + if (isRecord(m) && typeof m.name === 'string') mappings.set(m.name, m); + } + for (const job of jobs) { + if (!isRecord(job) || !isRecord(job.pull)) continue; + const name = job.pull.mapping; + if (typeof name !== 'string' || name.length === 0) continue; + const mapping = mappings.get(name); + if (!mapping) { + errors.push( + `Job '${String(job.name)}' pulls mapping '${name}' which is not defined in mappings. ` + + `Declare the mapping (its targetObject, fieldMapping and a connectorSource naming the ` + + `rest or openapi connector it pulls from) or correct the name.`, + ); + continue; + } + if (mapping.connectorSource === undefined) { + errors.push( + `Job '${String(job.name)}' pulls mapping '${name}', which declares no connectorSource — ` + + `there is nothing to pull. Add connectorSource: { connector, action } to the mapping, ` + + `naming the rest or openapi connector it reads from.`, + ); + } + } + return errors; +} + /** * Perform strict cross-reference validation on a parsed stack definition. * Returns an array of error messages (empty if valid). @@ -2851,6 +2896,9 @@ function validateCrossReferences( // Same placement, same reason: a global `operation: 'update'` action is // wrong whether or not the stack declares any objects. errors.push(...collectGlobalUpdateActionErrors(config)); + // Same placement: a job's `pull` names a MAPPING, resolved against the + // stack's mappings, and needs no object to resolve against. + errors.push(...collectJobPullMappingErrors(config)); if (objectNames.size === 0) return errors; diff --git a/packages/spec/src/system/job-pull-organization.test.ts b/packages/spec/src/system/job-pull-organization.test.ts new file mode 100644 index 00000000000..b9034f5ea22 --- /dev/null +++ b/packages/spec/src/system/job-pull-organization.test.ts @@ -0,0 +1,176 @@ +// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license. + +/** + * #20281 stage ③ — `JobSchema.pull` (ruling Q1-B) and `JobSchema.organization` + * (ruling Q2-O1), the contract half. + * + * - `pull: { mapping }` is a third run form, closed, and exclusive with `body` + * and `handler` — writing it beside either is refused at parse, located at + * `pull`. `body` + `handler` stays legal (the body wins), so the rule is + * pinned for both pairs and the old pair is pinned unchanged. + * - `defineStack` (and so `os validate`, whose config load runs it) refuses a + * `pull` naming a mapping the stack does not declare, or one with no + * `connectorSource`, with the cross-reference envelope. + * - `organization` is the scheduled flow's value shape, reused: a non-empty + * `sys_organization.id`. Its POSTURE rule is judged at bind, never here — + * pinned by the binder's suite (`packages/runtime`), not this one. + */ + +import { describe, expect, it } from 'vitest'; + +import { JobSchema, defineJob, type Job } from './job.zod'; +import { ScheduleOrganizationSchema } from '../automation/schedule-organization.zod'; +import { defineStack } from '../stack.zod'; + +const base = { name: 'orders_pull_hourly', schedule: { type: 'interval' as const, intervalMs: 60000 } }; +const body = { language: 'js' as const, source: "await ctx.api.object('task').find({});", capabilities: ['api.read' as const] }; +const pull = { mapping: 'orders_pull' }; + +const issueAt = (value: unknown, path: string) => { + const r = JobSchema.safeParse(value); + return r.success ? undefined : r.error.issues.find((i) => i.path.join('.') === path); +}; + +describe('JobSchema.pull — the declarative run form', () => { + it('accepts a job whose only run form is `pull` — it satisfies "something to run"', () => { + const r = JobSchema.safeParse({ ...base, pull }); + expect(r.success).toBe(true); + expect(r.success && r.data.pull).toEqual(pull); + expect(r.success && r.data.body).toBeUndefined(); + expect(r.success && r.data.handler).toBeUndefined(); + }); + + it('refuses `pull` beside `body`, located at `pull`', () => { + const issue = issueAt({ ...base, pull, body }, 'pull'); + expect(issue?.code).toBe('custom'); + expect(issue?.message).toContain('`pull`'); + }); + + it('refuses `pull` beside `handler`, located at `pull`', () => { + const issue = issueAt({ ...base, pull, handler: 'sweep' }, 'pull'); + expect(issue?.code).toBe('custom'); + }); + + it('keeps `body` + `handler` legal — the exclusion is the pull\'s, not a new rule on the old pair', () => { + expect(JobSchema.safeParse({ ...base, body, handler: 'sweep' }).success).toBe(true); + }); + + it('refuses a job with none of the three, located at `body`, naming all three keys', () => { + const issue = issueAt(base, 'body'); + expect(issue?.code).toBe('custom'); + for (const key of ['`body`', '`pull`', '`handler`']) expect(issue?.message).toContain(key); + }); + + it('is closed: an unrecognized key inside `pull` is refused, not stripped', () => { + const issue = issueAt({ ...base, pull: { ...pull, connector: 'orders_api' } }, 'pull'); + expect(issue?.code).toBe('unrecognized_keys'); + }); + + it('requires `mapping`, as a snake_case mapping name', () => { + expect(issueAt({ ...base, pull: {} }, 'pull.mapping')?.code).toBe('invalid_type'); + expect(issueAt({ ...base, pull: { mapping: 'OrdersPull' } }, 'pull.mapping')?.code).toBe('invalid_format'); + }); + + it('a bare top-level `mapping` is refused and pointed at `pull`', () => { + const r = JobSchema.safeParse({ ...base, mapping: 'orders_pull' }); + expect(r.success).toBe(false); + const issue = r.error!.issues.find((i) => i.code === 'unrecognized_keys'); + expect(issue?.message).toContain('`pull`'); + }); + + it('types `pull` on the authoring input', () => { + const job: Job = { ...base, pull }; + expect(defineJob(job).pull?.mapping).toBe('orders_pull'); + }); +}); + +describe('JobSchema.organization — the scheduled flow\'s value shape, reused', () => { + it('accepts a non-empty organization id on every run form', () => { + for (const form of [{ pull }, { body }, { handler: 'sweep' }]) { + const r = JobSchema.safeParse({ ...base, ...form, organization: 'org_a' }); + expect(r.success).toBe(true); + expect(r.success && r.data.organization).toBe('org_a'); + } + }); + + it('agrees with ScheduleOrganizationSchema on what counts as declared — one value shape, not a second', () => { + for (const value of ['org_a', 'x', '', 7, null]) { + const job = JobSchema.safeParse({ ...base, pull, organization: value }); + expect(job.success, JSON.stringify(value)).toBe(ScheduleOrganizationSchema.safeParse(value).success); + } + }); + + it('is optional at parse — whether it is required is the deployment\'s posture, read at bind', () => { + expect(JobSchema.safeParse({ ...base, pull }).success).toBe(true); + }); + + it('a near-miss spelling is refused at parse and pointed at `organization`', () => { + // `organization_id` has no entry of its own: the alias probe folds case and `_`. + for (const key of ['organizationId', 'organization_id', 'orgId', 'tenantId']) { + const r = JobSchema.safeParse({ ...base, pull, [key]: 'org_a' }); + expect(r.success, key).toBe(false); + const issue = r.error!.issues.find((i) => i.code === 'unrecognized_keys'); + expect(issue?.message, key).toContain('`organization`'); + } + }); +}); + +describe('defineStack refuses a job pull its stack cannot resolve — the cross-reference envelope', () => { + const manifest = { id: 'com.example.jobpull', name: 'job-pull-test', version: '1.0.0', type: 'app' as const }; + const order = { name: 'order', label: 'Order', fields: { external_id: { type: 'text' as const } } }; + const mapping = { + name: 'orders_pull', + targetObject: 'order', + fieldMapping: [{ source: 'id', target: 'external_id' }], + mode: 'upsert' as const, + upsertKey: ['external_id'], + connectorSource: { connector: 'orders_api', action: 'request' }, + }; + type Envelope = Error & { code?: string; status?: number; issues?: readonly string[] }; + const refusal = (extra: Record): Envelope | null => { + try { + defineStack({ manifest, objects: [order], ...extra } as unknown as Parameters[0]); + return null; + } catch (e) { + return e as Envelope; + } + }; + + it('accepts a pull naming a declared mapping with a connectorSource (control)', () => { + expect(refusal({ mappings: [mapping], jobs: [{ ...base, pull }] })).toBeNull(); + }); + + it('refuses a pull naming a mapping the stack does not declare', () => { + const refused = refusal({ mappings: [mapping], jobs: [{ ...base, pull: { mapping: 'orders_pul' } }] }); + expect(refused?.code).toBe('STACK_CROSS_REFERENCE_INVALID'); + expect(refused?.status).toBe(422); + expect(refused?.issues).toHaveLength(1); + expect(refused?.issues?.[0]).toContain("'orders_pull_hourly'"); + expect(refused?.issues?.[0]).toContain("'orders_pul'"); + }); + + it('refuses a pull naming a mapping with no connectorSource — there is nothing to pull', () => { + const { connectorSource: _dropped, ...importOnly } = mapping; + const refused = refusal({ mappings: [importOnly], jobs: [{ ...base, pull }] }); + expect(refused?.code).toBe('STACK_CROSS_REFERENCE_INVALID'); + expect(refused?.issues).toHaveLength(1); + expect(refused?.issues?.[0]).toContain('connectorSource'); + }); + + it('judges the reference on a stack that declares no object — it resolves against mappings alone', () => { + const refused = (() => { + try { + defineStack({ manifest, jobs: [{ ...base, pull }] } as unknown as Parameters[0]); + return null; + } catch (e) { + return e as Envelope; + } + })(); + expect(refused?.code).toBe('STACK_CROSS_REFERENCE_INVALID'); + }); + + it('judges a disabled job too — `enabled: false` turns a job off, it does not make a dangling name correct', () => { + const refused = refusal({ mappings: [mapping], jobs: [{ ...base, enabled: false, pull: { mapping: 'gone' } }] }); + expect(refused?.code).toBe('STACK_CROSS_REFERENCE_INVALID'); + }); +}); diff --git a/packages/spec/src/system/job.zod.ts b/packages/spec/src/system/job.zod.ts index 8736a9bdbbf..d84b4af1f8a 100644 --- a/packages/spec/src/system/job.zod.ts +++ b/packages/spec/src/system/job.zod.ts @@ -13,6 +13,8 @@ import { retiredKey } from '../shared/retired-key'; import { isValueDomainMember } from '../shared/value-domain.zod'; import { MetadataProtectionFields } from '../kernel/metadata-protection.zod'; import { ScriptBodySchema } from '../data/hook-body.zod'; +import { SnakeCaseIdentifierSchema } from '../shared/identifiers.zod'; +import { ScheduleOrganizationSchema } from '../automation/schedule-organization.zod'; import { bannedKeys, requiredOneOf } from '../shared/refinement-projection'; /** @@ -188,8 +190,47 @@ const JOB_TIMEOUT_RETIRED = */ const JOB_RUNS_NOTHING = 'A job needs something to run: declare `body` (a sandboxed JS body, `{ language: "js", source, ' - + 'capabilities }` — the preferred form) or `handler` (the name of a `defineStack({ functions })` ' - + 'entry — deprecated). This job declares neither, so no boot could ever schedule it.'; + + 'capabilities }` — the preferred form for code), `pull` (`{ mapping: "" }` — pull ' + + 'that mapping\'s `connectorSource` on this job\'s schedule, no code) or `handler` (the name of a ' + + '`defineStack({ functions })` entry — deprecated). This job declares none of them, so no boot ' + + 'could ever schedule it.'; + +/** + * A job's `pull` is a run form of its own, the declarative one (ruling Q1-B on + * the connector-sync card): the platform binds it and calls the automation + * service's pull, so there is no code beside it to run. `body` + `handler` + * stays legal (the body wins, the deprecated handler beside it is never run); + * `pull` + either is refused, because one of the two would silently never run. + */ +const JOB_PULL_EXCLUSIVE = + 'A job with `pull` runs that pull and nothing else: `pull` cannot be declared beside `body` or ' + + '`handler`, because the platform binds the pull itself and one of the two run forms would never ' + + 'run. Keep `pull` and delete the code, or keep the code (`body`) and delete `pull` — a body cannot ' + + 'call a pull; a pull on a schedule is this declaration.'; + +/** + * The organization a job runs as — its `body`'s `ctx.api`, the execution + * context its `handler` is handed, and its `pull`'s target read and writes + * alike. The value shape is the scheduled flow's + * ({@link ScheduleOrganizationSchema}, `sys_organization.id`), reused, and so + * is the deployment-posture rule that judges it: the maintainer's 2026-09-08 + * ruling on scheduled work across organizations, extended from scheduled flows + * to jobs by ruling Q2-O1 on the connector-sync card. That rule is + * `resolveScheduledWorkPolicy` (`@objectstack/types`), read by the job binder + * at BIND — never here: whether the key is required depends on the + * deployment's tenancy posture and its scheduled-work switch, and neither is + * knowable at authoring time (`automation/schedule-organization.zod.ts` says + * why an authoring-time rule for it must not exist). + */ +const JOB_ORGANIZATION_DESCRIPTION = + 'Organization id (sys_organization.id) this job runs as — its `body`\'s `ctx.api`, the execution ' + + 'context its `handler` is handed, and its `pull`\'s reads and writes alike, as a system run ' + + 'carrying that organization. A scheduled run has no session to inherit one from. Judged at bind ' + + 'by the posture rule scheduled flows use: required under the isolated tenancy posture (a job that ' + + 'declares none is not scheduled); optional under group (undeclared, the run carries no ' + + 'organization and a tenant-scoped write it makes is refused); not required under single (the ' + + 'install\'s one organization is resolved beneath each write). Where declared, it is the ' + + 'organization the run acts as on every posture.'; /** * The time limit of a body job has ONE spelling: the job's own `timeoutMs`. @@ -211,7 +252,16 @@ export const JobSchema = lazySchema(() => strictObject({ surface: 'this job', history: 'Until this shape was closed these were dropped silently — the item still registered, minus whatever the key was meant to configure.', - aliases: { cron: 'schedule', interval: 'schedule', fn: 'handler', function: 'handler', retry: 'retryPolicy', enabled_: 'enabled' }, + aliases: { + cron: 'schedule', interval: 'schedule', fn: 'handler', function: 'handler', retry: 'retryPolicy', enabled_: 'enabled', + // A pull written as a bare mapping name, or under the binding's own word. + mapping: 'pull', sync: 'pull', connectorSource: 'pull', + // The organization under the near-miss spellings the scheduled flow's + // refusal names — on this closed shape they are refused at parse, by name. + // One spelling per probe: the probe folds case and `_`, so `organization_id`, + // `org_id` and `tenant_id` are caught by the camelCase entries. + organizationId: 'organization', orgId: 'organization', tenantId: 'organization', tenant: 'organization', + }, guidance: { id: JOB_ID_RETIRED }, }, { // `id` removed in 17.0.0 (#4667) — see JOB_ID_RETIRED. `name` is the identity. @@ -228,10 +278,12 @@ export const JobSchema = lazySchema(() => strictObject({ * imports that module can bind. `body` is the form that travels with the * metadata itself, exactly as it did for hooks. * - * Optional since `body` exists, but not both absent: {@link JOB_RUNS_NOTHING}. - * When both are present `body` wins, as for hooks. + * Optional since `body` exists, but not all three of `body` / `handler` / + * `pull` absent: {@link JOB_RUNS_NOTHING}. When both `body` and `handler` are + * present `body` wins, as for hooks; beside `pull` it is refused + * ({@link JOB_PULL_EXCLUSIVE}). */ - handler: z.string().optional().describe('Handler function name (must match a key in `defineStack({ functions })`) — DEPRECATED, prefer `body`. When both are present `body` wins; a job must declare one of the two.'), + handler: z.string().optional().describe('Handler function name (must match a key in `defineStack({ functions })`) — DEPRECATED, prefer `body`. When both are present `body` wins; refused beside `pull`. A job must declare one of `body`, `handler` or `pull`.'), /** * Job Body (L2 sandboxed JS) — the hook body shape, reused by reference. * @@ -267,8 +319,49 @@ export const JobSchema = lazySchema(() => strictObject({ + 'Preferred over `handler`: when both are present `body` wins. ' + 'It runs in the QuickJS sandbox with no module scope (no imports, no helpers or constants from the surrounding file): it reaches data only through `ctx.api` under its declared `capabilities` (`api.read` / `api.write` / `api.transaction`) and logs through `ctx.log` (`log`); the in-process handler context (`ql`, `logger`, `bundle`) does not exist there. ' + "Its time limit is the job's `timeoutMs` (see there): long-running work declares a `timeoutMs` that covers it, or splits into bounded runs that each finish within it. " - + "Every door that brings an artifact in schedules a job's `body` — the boot, and `os package install` on install and on every restart — while a `handler` is code that travels only in the artifact's runtime module and runs only on a boot that loads it (a config, or `os start --artifact`); `os package install` therefore refuses an enabled job with no `body`.", + + "Every door that brings an artifact in schedules a job's `body` — the boot, and `os package install` on install and on every restart — while a `handler` is code that travels only in the artifact's runtime module and runs only on a boot that loads it (a config, or `os start --artifact`); `os package install` therefore refuses an enabled job with no `body` (a `pull` job excepted: it is data too). " + + 'Refused beside `pull`.', + ), + /** + * Job Pull — the declarative run form (ruling Q1-B on the connector-sync + * card): on each tick the platform pulls the named `mapping`'s + * `connectorSource` through the automation service + * (`IAutomationService.pullConnectorSource`) and writes the rows through + * the import runner. No code: the job says WHAT it does, so "which job pulls + * mapping M" is a metadata query, and the run's outcome is mapped once, by + * the binder — a refused pull is a `failed` run (retried per `retryPolicy`), + * a pull whose rows the import runner refused is `degraded`. + * + * A third run form beside `body` and `handler`, exclusive with both + * ({@link JOB_PULL_EXCLUSIVE}). It is data, so it travels with the artifact + * and binds on every door, as a `body` does. ⛔ Not a second code mechanism — + * a job carries code only as a `body` (the hook body shape, one binder) — + * and not a precedent for further task kinds: a second declarative member + * needs its own pull. + * + * `mapping` names a `mapping` the same stack declares, with a + * `connectorSource`: `defineStack` (and so `os validate`) refuses a name that + * resolves to neither. + */ + pull: strictObject({ + surface: 'this job’s pull', + history: + 'Declared closed from its first day: an unrecognized key is refused, never dropped.', + aliases: { name: 'mapping', mappingName: 'mapping', source: 'mapping', connectorSource: 'mapping', from: 'mapping' }, + }, { + mapping: SnakeCaseIdentifierSchema.describe( + 'Name of the `mapping` whose `connectorSource` this job pulls on its schedule — a mapping the same ' + + 'stack declares, with a `connectorSource` (refused at `defineStack` / `os validate` otherwise). The ' + + 'mapping says where the rows come from, the field map, the write mode and the match key; the job ' + + 'says when.', + ), + }).optional().describe( + 'Pull run form: on each run the platform pulls the named mapping\'s `connectorSource` (one action call, ' + + 'one response) and writes the rows through the import runner — no code. A refused pull records the run ' + + '`failed` (retried per `retryPolicy`); a pull whose rows the import runner refused records it ' + + '`degraded`. Data like `body`, so every door schedules it. Refused beside `body` or `handler`.', ), + organization: ScheduleOrganizationSchema.optional().describe(JOB_ORGANIZATION_DESCRIPTION), retryPolicy: RetryPolicySchema.optional().describe('Retry policy: failed runs (including timeouts) are retried with exponential backoff (delay = min(backoffMs * backoffMultiplier^(retry-1), maxRetryDelayMs), optionally jittered) up to maxRetries retries after the initial attempt. Omit the block for a single attempt; declaring it without `maxRetries` also means no retry since 17.0.0 — state a count to opt in.'), // Renamed from `timeout` (#14478): the unit (milliseconds) lived only in the // description while the sibling `retryPolicy.backoffMs` spells its own. @@ -292,7 +385,15 @@ export const JobSchema = lazySchema(() => strictObject({ // Declared through the closed projection list, so the published JSON Schema // states the rule (`anyOf` of one `required` per key) instead of being wider // than the parse. -}).refine(requiredOneOf(['body', 'handler']), { message: JOB_RUNS_NOTHING, path: ['body'] })); +}).refine(requiredOneOf(['body', 'handler', 'pull']), { message: JOB_RUNS_NOTHING, path: ['body'] }) +// "`pull` excludes `body` and `handler`" has no arm in the closed projection +// list, so the published JSON Schema stays wider than this rule; the site is +// recorded in `dropped-refinements.baseline.json`, beside every other +// mutual-exclusion refinement in this package. + .refine((job) => job.pull === undefined || (job.body === undefined && job.handler === undefined), { + message: JOB_PULL_EXCLUSIVE, + path: ['pull'], + })); export type Job = z.input; /** Post-parse shape of {@link Job} — defaults applied, transforms run (ADR-0122). */