Press n or j to go to the next uncovered block, b, p or k for the previous block.
| 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 | 28x 28x 28x 12x 12x 5x 3x 3x 1x 2x 2x 1x 1x 1x | /**
* @module
*
* @author Bayu Dwiyan Satria
* @version 1.0.0
* @since 1.3.0
*/
import { Binding } from '@/core/bindings/Binding'
import { lazySettings } from '@/utils/lazySettings'
import { PipelineSettings } from '@/types/PipelineSettings'
import type { Pipeline, PipelineRecord } from 'cloudflare:pipelines'
import type { CloudflareEnv } from '@/types/CloudflareEnv'
/**
* Accessor for a Pipelines binding.
*
* @remarks
* A pipeline takes structured records and lands them somewhere durable — R2,
* as files, batched and partitioned by rules configured on the pipeline rather
* than here. The Worker sends records and stops thinking about them.
*
* That last part is the whole point, and it is what separates this from
* {@link R2Binding}. Writing to R2 directly means a producer owns buffering:
* how many records make a file, how long to wait for a slow hour, what to do
* with a partial batch when the isolate is evicted. Ingestion that arrives in
* bursts — a large batch produced inside one cron window — are exactly where
* that goes wrong, and getting it wrong loses data quietly.
*
* It is also not {@link QueueBinding}, though both take a batch and return
* quickly. A queue exists so another Worker can act on each message; a pipeline
* exists so nothing has to. If the records are going to be read as files later
* rather than processed one at a time, the queue's consumer is a step that only
* exists to write them down.
*
* ## Two ways to send, and the safe one is the point
*
* {@link send} throws, like every other accessor here. {@link sendSafe} logs
* and returns `false`, mirroring {@link AnalyticsBinding}'s `writeSafe` — and
* for the same reason. A pipeline's usual job is a second, independent sink
* beside the one doing the real work: archiving the raw record a provider
* returned while the coerced one goes on to storage. A sink like that failing
* the run that produced the records inverts its purpose, because the archive
* exists to explain a bad run and would instead be causing one.
*
* Which to call follows from whether the record has anywhere else to be. An
* audit trail that is the only copy should throw.
*
* ## Why there is no `PipelineService`, and no transformation support
*
* No kernel capability to implement, as with {@link WorkflowBinding} — and a
* pipeline is closer to plumbing than to a capability in any case.
*
* Transformation is a separate matter. A pipeline can run records through a
* `PipelineTransformationEntrypoint` before landing them, which — like a
* Workflow's `WorkflowEntrypoint` — is a class the consuming Worker exports and
* names in its own `wrangler.json`. A base class that must be exported from
* someone else's entry has nothing an adapter can add, so this package wraps
* the producer side only.
*
* @typeParam T The record shape this pipeline carries.
*
* @example
* ```ts
* const archive = new PipelineBinding<RawRecord>()
*
* // Archival must not fail the ingestion run that produced the records.
* await archive.sendSafe(env, rejected.map(record => ({ ...record, executionId, retrievedAt })))
* ```
*
* @class
*
* @author Bayu Dwiyan Satria
* @version 1.0.0
* @since 1.3.0
*/
export class PipelineBinding<
T extends PipelineRecord = PipelineRecord,
E extends CloudflareEnv = CloudflareEnv
> extends Binding<Pipeline<T>, E> {
/**
* Constructs a PipelineBinding instance.
*
* @param settings Resolved settings. Defaults to `resolve('pipeline')`, so callers
* normally construct this with no arguments at all.
*/
constructor(settings?: PipelineSettings) {
const settle = lazySettings<PipelineSettings>('pipeline', settings)
super(() => settle().binding)
}
/**
* Sends records to the pipeline.
*
* @remarks
* Resolves when the pipeline has accepted the records, not when they have
* been written — the batching that makes a pipeline worth using happens after
* this returns, so a successful send is a promise to land the records rather
* than evidence that it happened.
*
* An empty array is sent as given rather than short-circuited, because
* whether an empty batch is worth a call is the pipeline's business and a
* caller that filtered everything out has already decided to send.
*
* @param env The Worker environment.
* @param records The records to send.
* @throws {@link MissingBindingError} When the pipeline binding is absent.
*/
public async send(env: E, records: T[]): Promise<void> {
await this.resolve(env).send(records)
}
/**
* Sends records without letting a failure reach the caller.
*
* @remarks
* The method an archival sink should call. A missing binding, a rejected
* batch, an unreachable pipeline — each is logged and reported through the
* return value, so the work that produced the records finishes either way.
*
* Reach for {@link send} instead when the pipeline is the only copy: silence
* is the right answer for a second sink and the wrong one for the only sink.
*
* @param env The Worker environment.
* @param records The records to send.
* @returns `true` when the pipeline accepted the records.
*/
public async sendSafe(env: E, records: T[]): Promise<boolean> {
const pipeline = this.tryResolve(env)
if (pipeline === null) {
return false
}
try {
await pipeline.send(records)
return true
} catch (e) {
console.error('[Pipeline Error]', e)
return false
}
}
}
|