AK-21/Graphite-Industrial-Intelligence
0
1import { __commonJS, __toESM, require_defineProperty, require_objectSpread2 } from "./objectSpread2-BvkFp-_Y.mjs";2import { createChain, splitLink } from "./splitLink-B7Cuf2c_.mjs";3import { TRPCClientError, isTRPCClientError } from "./TRPCClientError-apv8gw59.mjs";4import { fetchHTTPResponse, getBody, getFetch, getUrl, resolveHTTPLinkOptions } from "./httpUtils-BNq9QC3d.mjs";5import { httpLink, isFormData, isNonJsonSerializable, isOctetType } from "./httpLink-oiU8eqFi.mjs";6import { abortSignalToPromise, allAbortSignals, dataLoader, httpBatchLink, raceAbortSignals } from "./httpBatchLink-CaWjh1oW.mjs";7import { getTransformer } from "./unstable-internals-Bg7n9BBj.mjs";8import { loggerLink } from "./loggerLink-ineCN1PO.mjs";9import { createWSClient, jsonEncoder, resultOf, wsLink } from "./wsLink-DSf4KOdW.mjs";10import { behaviorSubject, observable, observableToPromise, share } from "@trpc/server/observable";11import { callProcedure, createFlatProxy, createRecursiveProxy, isAbortError, isAsyncIterable, iteratorResource, jsonlStreamConsumer, makeResource, retryableRpcCodes, run, sseStreamConsumer } from "@trpc/server/unstable-core-do-not-import";12import { getTRPCErrorFromUnknown, getTRPCErrorShape, isTrackedEnvelope } from "@trpc/server";13import { TRPC_ERROR_CODES_BY_KEY } from "@trpc/server/rpc";14 15//#region src/internals/TRPCUntypedClient.ts16var import_defineProperty = __toESM(require_defineProperty(), 1);17var import_objectSpread2$4 = __toESM(require_objectSpread2(), 1);18var TRPCUntypedClient = class {19 constructor(opts) {20 (0, import_defineProperty.default)(this, "links", void 0);21 (0, import_defineProperty.default)(this, "runtime", void 0);22 (0, import_defineProperty.default)(this, "requestId", void 0);23 this.requestId = 0;24 this.runtime = {};25 this.links = opts.links.map((link) => link(this.runtime));26 }27 $request(opts) {28 var _opts$context;29 const chain$ = createChain({30 links: this.links,31 op: (0, import_objectSpread2$4.default)((0, import_objectSpread2$4.default)({}, opts), {}, {32 context: (_opts$context = opts.context) !== null && _opts$context !== void 0 ? _opts$context : {},33 id: ++this.requestId34 })35 });36 return chain$.pipe(share());37 }38 async requestAsPromise(opts) {39 var _this = this;40 try {41 const req$ = _this.$request(opts);42 const envelope = await observableToPromise(req$);43 const data = envelope.result.data;44 return data;45 } catch (err) {46 throw TRPCClientError.from(err);47 }48 }49 query(path, input, opts) {50 return this.requestAsPromise({51 type: "query",52 path,53 input,54 context: opts === null || opts === void 0 ? void 0 : opts.context,55 signal: opts === null || opts === void 0 ? void 0 : opts.signal56 });57 }58 mutation(path, input, opts) {59 return this.requestAsPromise({60 type: "mutation",61 path,62 input,63 context: opts === null || opts === void 0 ? void 0 : opts.context,64 signal: opts === null || opts === void 0 ? void 0 : opts.signal65 });66 }67 subscription(path, input, opts) {68 const observable$ = this.$request({69 type: "subscription",70 path,71 input,72 context: opts.context,73 signal: opts.signal74 });75 return observable$.subscribe({76 next(envelope) {77 switch (envelope.result.type) {78 case "state": {79 var _opts$onConnectionSta;80 (_opts$onConnectionSta = opts.onConnectionStateChange) === null || _opts$onConnectionSta === void 0 || _opts$onConnectionSta.call(opts, envelope.result);81 break;82 }83 case "started": {84 var _opts$onStarted;85 (_opts$onStarted = opts.onStarted) === null || _opts$onStarted === void 0 || _opts$onStarted.call(opts, { context: envelope.context });86 break;87 }88 case "stopped": {89 var _opts$onStopped;90 (_opts$onStopped = opts.onStopped) === null || _opts$onStopped === void 0 || _opts$onStopped.call(opts);91 break;92 }93 case "data":94 case void 0: {95 var _opts$onData;96 (_opts$onData = opts.onData) === null || _opts$onData === void 0 || _opts$onData.call(opts, envelope.result.data);97 break;98 }99 }100 },101 error(err) {102 var _opts$onError;103 (_opts$onError = opts.onError) === null || _opts$onError === void 0 || _opts$onError.call(opts, err);104 },105 complete() {106 var _opts$onComplete;107 (_opts$onComplete = opts.onComplete) === null || _opts$onComplete === void 0 || _opts$onComplete.call(opts);108 }109 });110 }111};112 113//#endregion114//#region src/createTRPCUntypedClient.ts115function createTRPCUntypedClient(opts) {116 return new TRPCUntypedClient(opts);117}118 119//#endregion120//#region src/createTRPCClient.ts121const untypedClientSymbol = Symbol.for("trpc_untypedClient");122const clientCallTypeMap = {123 query: "query",124 mutate: "mutation",125 subscribe: "subscription"126};127/** @internal */128const clientCallTypeToProcedureType = (clientCallType) => {129 return clientCallTypeMap[clientCallType];130};131/**132* @internal133*/134function createTRPCClientProxy(client) {135 const proxy = createRecursiveProxy(({ path, args }) => {136 const pathCopy = [...path];137 const procedureType = clientCallTypeToProcedureType(pathCopy.pop());138 const fullPath = pathCopy.join(".");139 return client[procedureType](fullPath, ...args);140 });141 return createFlatProxy((key) => {142 if (key === untypedClientSymbol) return client;143 return proxy[key];144 });145}146function createTRPCClient(opts) {147 const client = new TRPCUntypedClient(opts);148 const proxy = createTRPCClientProxy(client);149 return proxy;150}151/**152* Get an untyped client from a proxy client153* @internal154*/155function getUntypedClient(client) {156 return client[untypedClientSymbol];157}158 159//#endregion160//#region src/links/httpBatchStreamLink.ts161var import_objectSpread2$3 = __toESM(require_objectSpread2(), 1);162/**163* @see https://trpc.io/docs/client/links/httpBatchStreamLink164*/165function httpBatchStreamLink(opts) {166 var _opts$maxURLLength, _opts$maxItems;167 const resolvedOpts = resolveHTTPLinkOptions(opts);168 const maxURLLength = (_opts$maxURLLength = opts.maxURLLength) !== null && _opts$maxURLLength !== void 0 ? _opts$maxURLLength : Infinity;169 const maxItems = (_opts$maxItems = opts.maxItems) !== null && _opts$maxItems !== void 0 ? _opts$maxItems : Infinity;170 return () => {171 const batchLoader = (type) => {172 return {173 validate(batchOps) {174 if (maxURLLength === Infinity && maxItems === Infinity) return true;175 if (batchOps.length > maxItems) return false;176 const path = batchOps.map((op) => op.path).join(",");177 const inputs = batchOps.map((op) => op.input);178 const url = getUrl((0, import_objectSpread2$3.default)((0, import_objectSpread2$3.default)({}, resolvedOpts), {}, {179 type,180 path,181 inputs,182 signal: null183 }));184 return url.length <= maxURLLength;185 },186 async fetch(batchOps) {187 var _opts$streamHeader;188 const path = batchOps.map((op) => op.path).join(",");189 const inputs = batchOps.map((op) => op.input);190 const batchSignals = allAbortSignals(...batchOps.map((op) => op.signal));191 const abortController = new AbortController();192 const responsePromise = fetchHTTPResponse((0, import_objectSpread2$3.default)((0, import_objectSpread2$3.default)({}, resolvedOpts), {}, {193 signal: raceAbortSignals(batchSignals, abortController.signal),194 type,195 contentTypeHeader: "application/json",196 trpcAcceptHeader: "application/jsonl",197 trpcAcceptHeaderKey: (_opts$streamHeader = opts.streamHeader) !== null && _opts$streamHeader !== void 0 ? _opts$streamHeader : "trpc-accept",198 getUrl,199 getBody,200 inputs,201 path,202 headers() {203 if (!opts.headers) return {};204 if (typeof opts.headers === "function") return opts.headers({ opList: batchOps });205 return opts.headers;206 }207 }));208 const res = await responsePromise;209 if (!res.ok) {210 const json = await res.json();211 if ("error" in json) json.error = resolvedOpts.transformer.output.deserialize(json.error);212 return batchOps.map(() => Promise.resolve({213 json,214 meta: { response: res }215 }));216 }217 const [head] = await jsonlStreamConsumer({218 from: res.body,219 deserialize: (data) => resolvedOpts.transformer.output.deserialize(data),220 formatError(opts$1) {221 const error = opts$1.error;222 return TRPCClientError.from({ error });223 },224 abortController225 });226 const promises = Object.keys(batchOps).map(async (key) => {227 let json = await Promise.resolve(head[key]);228 if ("result" in json) {229 /**230 * Not very pretty, but we need to unwrap nested data as promises231 * Our stream producer will only resolve top-level async values or async values that are directly nested in another async value232 */233 const result = await Promise.resolve(json.result);234 json = { result: { data: await Promise.resolve(result.data) } };235 }236 return {237 json,238 meta: { response: res }239 };240 });241 return promises;242 }243 };244 };245 const query = dataLoader(batchLoader("query"));246 const mutation = dataLoader(batchLoader("mutation"));247 const loaders = {248 query,249 mutation250 };251 return ({ op }) => {252 return observable((observer) => {253 /* istanbul ignore if -- @preserve */254 if (op.type === "subscription") throw new Error("Subscriptions are unsupported by `httpBatchStreamLink` - use `httpSubscriptionLink` or `wsLink`");255 const loader = loaders[op.type];256 const promise = loader.load(op);257 let _res = void 0;258 promise.then((res) => {259 _res = res;260 if ("error" in res.json) {261 observer.error(TRPCClientError.from(res.json, { meta: res.meta }));262 return;263 } else if ("result" in res.json) {264 observer.next({265 context: res.meta,266 result: res.json.result267 });268 observer.complete();269 return;270 }271 observer.complete();272 }).catch((err) => {273 observer.error(TRPCClientError.from(err, { meta: _res === null || _res === void 0 ? void 0 : _res.meta }));274 });275 return () => {};276 });277 };278 };279}280/**281* @deprecated use {@link httpBatchStreamLink} instead282*/283const unstable_httpBatchStreamLink = httpBatchStreamLink;284 285//#endregion286//#region src/internals/inputWithTrackedEventId.ts287var import_objectSpread2$2 = __toESM(require_objectSpread2(), 1);288function inputWithTrackedEventId(input, lastEventId) {289 if (!lastEventId) return input;290 if (input != null && typeof input !== "object") return input;291 return (0, import_objectSpread2$2.default)((0, import_objectSpread2$2.default)({}, input !== null && input !== void 0 ? input : {}), {}, { lastEventId });292}293 294//#endregion295//#region ../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/asyncIterator.js296var require_asyncIterator = __commonJS({ "../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/asyncIterator.js"(exports, module) {297 function _asyncIterator$1(r) {298 var n, t, o, e = 2;299 for ("undefined" != typeof Symbol && (t = Symbol.asyncIterator, o = Symbol.iterator); e--;) {300 if (t && null != (n = r[t])) return n.call(r);301 if (o && null != (n = r[o])) return new AsyncFromSyncIterator(n.call(r));302 t = "@@asyncIterator", o = "@@iterator";303 }304 throw new TypeError("Object is not async iterable");305 }306 function AsyncFromSyncIterator(r) {307 function AsyncFromSyncIteratorContinuation(r$1) {308 if (Object(r$1) !== r$1) return Promise.reject(new TypeError(r$1 + " is not an object."));309 var n = r$1.done;310 return Promise.resolve(r$1.value).then(function(r$2) {311 return {312 value: r$2,313 done: n314 };315 });316 }317 return AsyncFromSyncIterator = function AsyncFromSyncIterator$1(r$1) {318 this.s = r$1, this.n = r$1.next;319 }, AsyncFromSyncIterator.prototype = {320 s: null,321 n: null,322 next: function next() {323 return AsyncFromSyncIteratorContinuation(this.n.apply(this.s, arguments));324 },325 "return": function _return(r$1) {326 var n = this.s["return"];327 return void 0 === n ? Promise.resolve({328 value: r$1,329 done: !0330 }) : AsyncFromSyncIteratorContinuation(n.apply(this.s, arguments));331 },332 "throw": function _throw(r$1) {333 var n = this.s["return"];334 return void 0 === n ? Promise.reject(r$1) : AsyncFromSyncIteratorContinuation(n.apply(this.s, arguments));335 }336 }, new AsyncFromSyncIterator(r);337 }338 module.exports = _asyncIterator$1, module.exports.__esModule = true, module.exports["default"] = module.exports;339} });340 341//#endregion342//#region src/links/httpSubscriptionLink.ts343var import_asyncIterator = __toESM(require_asyncIterator(), 1);344async function urlWithConnectionParams(opts) {345 let url = await resultOf(opts.url);346 if (opts.connectionParams) {347 const params = await resultOf(opts.connectionParams);348 const prefix = url.includes("?") ? "&" : "?";349 url += prefix + "connectionParams=" + encodeURIComponent(JSON.stringify(params));350 }351 return url;352}353/**354* @see https://trpc.io/docs/client/links/httpSubscriptionLink355*/356function httpSubscriptionLink(opts) {357 const transformer = getTransformer(opts.transformer);358 return () => {359 return ({ op }) => {360 return observable((observer) => {361 var _opts$EventSource;362 const { type, path, input } = op;363 /* istanbul ignore if -- @preserve */364 if (type !== "subscription") throw new Error("httpSubscriptionLink only supports subscriptions");365 let lastEventId = void 0;366 const ac = new AbortController();367 const signal = raceAbortSignals(op.signal, ac.signal);368 const eventSourceStream = sseStreamConsumer({369 url: async () => getUrl({370 transformer,371 url: await urlWithConnectionParams(opts),372 input: inputWithTrackedEventId(input, lastEventId),373 path,374 type,375 signal: null376 }),377 init: () => resultOf(opts.eventSourceOptions, { op }),378 signal,379 deserialize: (data) => transformer.output.deserialize(data),380 EventSource: (_opts$EventSource = opts.EventSource) !== null && _opts$EventSource !== void 0 ? _opts$EventSource : globalThis.EventSource381 });382 const connectionState = behaviorSubject({383 type: "state",384 state: "connecting",385 error: null386 });387 const connectionSub = connectionState.subscribe({ next(state) {388 observer.next({ result: state });389 } });390 run(async () => {391 var _iteratorAbruptCompletion = false;392 var _didIteratorError = false;393 var _iteratorError;394 try {395 for (var _iterator = (0, import_asyncIterator.default)(eventSourceStream), _step; _iteratorAbruptCompletion = !(_step = await _iterator.next()).done; _iteratorAbruptCompletion = false) {396 const chunk = _step.value;397 switch (chunk.type) {398 case "ping": break;399 case "data":400 const chunkData = chunk.data;401 let result;402 if (chunkData.id) {403 lastEventId = chunkData.id;404 result = {405 id: chunkData.id,406 data: chunkData407 };408 } else result = { data: chunkData.data };409 observer.next({410 result,411 context: { eventSource: chunk.eventSource }412 });413 break;414 case "connected": {415 observer.next({416 result: { type: "started" },417 context: { eventSource: chunk.eventSource }418 });419 connectionState.next({420 type: "state",421 state: "pending",422 error: null423 });424 break;425 }426 case "serialized-error": {427 const error = TRPCClientError.from({ error: chunk.error });428 if (retryableRpcCodes.includes(chunk.error.code)) {429 connectionState.next({430 type: "state",431 state: "connecting",432 error433 });434 break;435 }436 throw error;437 }438 case "connecting": {439 const lastState = connectionState.get();440 const error = chunk.event && TRPCClientError.from(chunk.event);441 if (!error && lastState.state === "connecting") break;442 connectionState.next({443 type: "state",444 state: "connecting",445 error446 });447 break;448 }449 case "timeout": connectionState.next({450 type: "state",451 state: "connecting",452 error: new TRPCClientError(`Timeout of ${chunk.ms}ms reached while waiting for a response`)453 });454 }455 }456 } catch (err) {457 _didIteratorError = true;458 _iteratorError = err;459 } finally {460 try {461 if (_iteratorAbruptCompletion && _iterator.return != null) await _iterator.return();462 } finally {463 if (_didIteratorError) throw _iteratorError;464 }465 }466 observer.next({ result: { type: "stopped" } });467 connectionState.next({468 type: "state",469 state: "idle",470 error: null471 });472 observer.complete();473 }).catch((error) => {474 observer.error(TRPCClientError.from(error));475 });476 return () => {477 observer.complete();478 ac.abort();479 connectionSub.unsubscribe();480 };481 });482 };483 };484}485/**486* @deprecated use {@link httpSubscriptionLink} instead487*/488const unstable_httpSubscriptionLink = httpSubscriptionLink;489 490//#endregion491//#region src/links/retryLink.ts492var import_objectSpread2$1 = __toESM(require_objectSpread2(), 1);493/**494* @see https://trpc.io/docs/v11/client/links/retryLink495*/496function retryLink(opts) {497 return () => {498 return (callOpts) => {499 return observable((observer) => {500 let next$;501 let callNextTimeout = void 0;502 let lastEventId = void 0;503 attempt(1);504 function opWithLastEventId() {505 const op = callOpts.op;506 if (!lastEventId) return op;507 return (0, import_objectSpread2$1.default)((0, import_objectSpread2$1.default)({}, op), {}, { input: inputWithTrackedEventId(op.input, lastEventId) });508 }509 function attempt(attempts) {510 const op = opWithLastEventId();511 next$ = callOpts.next(op).subscribe({512 error(error) {513 var _opts$retryDelayMs, _opts$retryDelayMs2;514 const shouldRetry = opts.retry({515 op,516 attempts,517 error518 });519 if (!shouldRetry) {520 observer.error(error);521 return;522 }523 const delayMs = (_opts$retryDelayMs = (_opts$retryDelayMs2 = opts.retryDelayMs) === null || _opts$retryDelayMs2 === void 0 ? void 0 : _opts$retryDelayMs2.call(opts, attempts)) !== null && _opts$retryDelayMs !== void 0 ? _opts$retryDelayMs : 0;524 if (delayMs <= 0) {525 attempt(attempts + 1);526 return;527 }528 callNextTimeout = setTimeout(() => attempt(attempts + 1), delayMs);529 },530 next(envelope) {531 if ((!envelope.result.type || envelope.result.type === "data") && envelope.result.id) lastEventId = envelope.result.id;532 observer.next(envelope);533 },534 complete() {535 observer.complete();536 }537 });538 }539 return () => {540 next$.unsubscribe();541 clearTimeout(callNextTimeout);542 };543 });544 };545 };546}547 548//#endregion549//#region ../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/usingCtx.js550var require_usingCtx = __commonJS({ "../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/usingCtx.js"(exports, module) {551 function _usingCtx() {552 var r = "function" == typeof SuppressedError ? SuppressedError : function(r$1, e$1) {553 var n$1 = Error();554 return n$1.name = "SuppressedError", n$1.error = r$1, n$1.suppressed = e$1, n$1;555 }, e = {}, n = [];556 function using(r$1, e$1) {557 if (null != e$1) {558 if (Object(e$1) !== e$1) throw new TypeError("using declarations can only be used with objects, functions, null, or undefined.");559 if (r$1) var o = e$1[Symbol.asyncDispose || Symbol["for"]("Symbol.asyncDispose")];560 if (void 0 === o && (o = e$1[Symbol.dispose || Symbol["for"]("Symbol.dispose")], r$1)) var t = o;561 if ("function" != typeof o) throw new TypeError("Object is not disposable.");562 t && (o = function o$1() {563 try {564 t.call(e$1);565 } catch (r$2) {566 return Promise.reject(r$2);567 }568 }), n.push({569 v: e$1,570 d: o,571 a: r$1572 });573 } else r$1 && n.push({574 d: e$1,575 a: r$1576 });577 return e$1;578 }579 return {580 e,581 u: using.bind(null, !1),582 a: using.bind(null, !0),583 d: function d() {584 var o, t = this.e, s = 0;585 function next() {586 for (; o = n.pop();) try {587 if (!o.a && 1 === s) return s = 0, n.push(o), Promise.resolve().then(next);588 if (o.d) {589 var r$1 = o.d.call(o.v);590 if (o.a) return s |= 2, Promise.resolve(r$1).then(next, err);591 } else s |= 1;592 } catch (r$2) {593 return err(r$2);594 }595 if (1 === s) return t !== e ? Promise.reject(t) : Promise.resolve();596 if (t !== e) throw t;597 }598 function err(n$1) {599 return t = t !== e ? new r(n$1, t) : n$1, next();600 }601 return next();602 }603 };604 }605 module.exports = _usingCtx, module.exports.__esModule = true, module.exports["default"] = module.exports;606} });607 608//#endregion609//#region ../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/OverloadYield.js610var require_OverloadYield = __commonJS({ "../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/OverloadYield.js"(exports, module) {611 function _OverloadYield(e, d) {612 this.v = e, this.k = d;613 }614 module.exports = _OverloadYield, module.exports.__esModule = true, module.exports["default"] = module.exports;615} });616 617//#endregion618//#region ../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/awaitAsyncGenerator.js619var require_awaitAsyncGenerator = __commonJS({ "../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/awaitAsyncGenerator.js"(exports, module) {620 var OverloadYield$1 = require_OverloadYield();621 function _awaitAsyncGenerator$1(e) {622 return new OverloadYield$1(e, 0);623 }624 module.exports = _awaitAsyncGenerator$1, module.exports.__esModule = true, module.exports["default"] = module.exports;625} });626 627//#endregion628//#region ../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/wrapAsyncGenerator.js629var require_wrapAsyncGenerator = __commonJS({ "../../node_modules/.pnpm/@oxc-project+runtime@0.72.2/node_modules/@oxc-project/runtime/src/helpers/wrapAsyncGenerator.js"(exports, module) {630 var OverloadYield = require_OverloadYield();631 function _wrapAsyncGenerator$1(e) {632 return function() {633 return new AsyncGenerator(e.apply(this, arguments));634 };635 }636 function AsyncGenerator(e) {637 var r, t;638 function resume(r$1, t$1) {639 try {640 var n = e[r$1](t$1), o = n.value, u = o instanceof OverloadYield;641 Promise.resolve(u ? o.v : o).then(function(t$2) {642 if (u) {643 var i = "return" === r$1 ? "return" : "next";644 if (!o.k || t$2.done) return resume(i, t$2);645 t$2 = e[i](t$2).value;646 }647 settle(n.done ? "return" : "normal", t$2);648 }, function(e$1) {649 resume("throw", e$1);650 });651 } catch (e$1) {652 settle("throw", e$1);653 }654 }655 function settle(e$1, n) {656 switch (e$1) {657 case "return":658 r.resolve({659 value: n,660 done: !0661 });662 break;663 case "throw":664 r.reject(n);665 break;666 default: r.resolve({667 value: n,668 done: !1669 });670 }671 (r = r.next) ? resume(r.key, r.arg) : t = null;672 }673 this._invoke = function(e$1, n) {674 return new Promise(function(o, u) {675 var i = {676 key: e$1,677 arg: n,678 resolve: o,679 reject: u,680 next: null681 };682 t ? t = t.next = i : (r = t = i, resume(e$1, n));683 });684 }, "function" != typeof e["return"] && (this["return"] = void 0);685 }686 AsyncGenerator.prototype["function" == typeof Symbol && Symbol.asyncIterator || "@@asyncIterator"] = function() {687 return this;688 }, AsyncGenerator.prototype.next = function(e) {689 return this._invoke("next", e);690 }, AsyncGenerator.prototype["throw"] = function(e) {691 return this._invoke("throw", e);692 }, AsyncGenerator.prototype["return"] = function(e) {693 return this._invoke("return", e);694 };695 module.exports = _wrapAsyncGenerator$1, module.exports.__esModule = true, module.exports["default"] = module.exports;696} });697 698//#endregion699//#region src/links/localLink.ts700var import_usingCtx = __toESM(require_usingCtx(), 1);701var import_awaitAsyncGenerator = __toESM(require_awaitAsyncGenerator(), 1);702var import_wrapAsyncGenerator = __toESM(require_wrapAsyncGenerator(), 1);703var import_objectSpread2 = __toESM(require_objectSpread2(), 1);704/**705* localLink is a terminating link that allows you to make tRPC procedure calls directly in your application without going through HTTP.706*707* @see https://trpc.io/docs/links/localLink708*/709function unstable_localLink(opts) {710 const transformer = getTransformer(opts.transformer);711 const transformChunk = (chunk) => {712 if (opts.transformer) return chunk;713 if (chunk === void 0) return chunk;714 const serialized = JSON.stringify(transformer.input.serialize(chunk));715 const deserialized = JSON.parse(transformer.output.deserialize(serialized));716 return deserialized;717 };718 return () => ({ op }) => observable((observer) => {719 let ctx = void 0;720 const ac = new AbortController();721 const signal = raceAbortSignals(op.signal, ac.signal);722 const signalPromise = abortSignalToPromise(signal);723 signalPromise.catch(() => {});724 let input = op.input;725 async function runProcedure(newInput) {726 input = newInput;727 ctx = await opts.createContext();728 return callProcedure({729 router: opts.router,730 path: op.path,731 getRawInput: async () => newInput,732 ctx,733 type: op.type,734 signal,735 batchIndex: 0736 });737 }738 function onErrorCallback(cause) {739 var _opts$onError;740 if (isAbortError(cause)) return;741 (_opts$onError = opts.onError) === null || _opts$onError === void 0 || _opts$onError.call(opts, {742 error: getTRPCErrorFromUnknown(cause),743 type: op.type,744 path: op.path,745 input,746 ctx747 });748 }749 function coerceToTRPCClientError(cause) {750 if (isTRPCClientError(cause)) return cause;751 const error = getTRPCErrorFromUnknown(cause);752 const shape = getTRPCErrorShape({753 config: opts.router._def._config,754 ctx,755 error,756 input,757 path: op.path,758 type: op.type759 });760 return TRPCClientError.from({ error: transformChunk(shape) }, { cause: cause instanceof Error ? cause : void 0 });761 }762 run(async () => {763 switch (op.type) {764 case "query":765 case "mutation": {766 const result = await runProcedure(op.input);767 if (!isAsyncIterable(result)) {768 observer.next({ result: { data: transformChunk(result) } });769 observer.complete();770 break;771 }772 observer.next({ result: { data: (0, import_wrapAsyncGenerator.default)(function* () {773 try {774 var _usingCtx$1 = (0, import_usingCtx.default)();775 const iterator = _usingCtx$1.a(iteratorResource(result));776 const _finally = _usingCtx$1.u(makeResource({}, () => {777 observer.complete();778 }));779 try {780 while (true) {781 const res = yield (0, import_awaitAsyncGenerator.default)(Promise.race([iterator.next(), signalPromise]));782 if (res.done) return transformChunk(res.value);783 yield transformChunk(res.value);784 }785 } catch (cause) {786 onErrorCallback(cause);787 throw coerceToTRPCClientError(cause);788 }789 } catch (_) {790 _usingCtx$1.e = _;791 } finally {792 yield (0, import_awaitAsyncGenerator.default)(_usingCtx$1.d());793 }794 })() } });795 break;796 }797 case "subscription": try {798 var _usingCtx3 = (0, import_usingCtx.default)();799 const connectionState = behaviorSubject({800 type: "state",801 state: "connecting",802 error: null803 });804 const connectionSub = connectionState.subscribe({ next(state) {805 observer.next({ result: state });806 } });807 let lastEventId = void 0;808 const _finally = _usingCtx3.u(makeResource({}, async () => {809 observer.complete();810 connectionState.next({811 type: "state",812 state: "idle",813 error: null814 });815 connectionSub.unsubscribe();816 }));817 while (true) try {818 var _usingCtx4 = (0, import_usingCtx.default)();819 const result = await runProcedure(inputWithTrackedEventId(op.input, lastEventId));820 if (!isAsyncIterable(result)) throw new Error("Expected an async iterable");821 const iterator = _usingCtx4.a(iteratorResource(result));822 observer.next({ result: { type: "started" } });823 connectionState.next({824 type: "state",825 state: "pending",826 error: null827 });828 while (true) {829 let res;830 try {831 res = await Promise.race([iterator.next(), signalPromise]);832 } catch (cause) {833 if (isAbortError(cause)) return;834 const error = getTRPCErrorFromUnknown(cause);835 if (!retryableRpcCodes.includes(TRPC_ERROR_CODES_BY_KEY[error.code])) throw coerceToTRPCClientError(error);836 onErrorCallback(error);837 connectionState.next({838 type: "state",839 state: "connecting",840 error: coerceToTRPCClientError(error)841 });842 break;843 }844 if (res.done) return;845 let chunk;846 if (isTrackedEnvelope(res.value)) {847 lastEventId = res.value[0];848 chunk = {849 id: res.value[0],850 data: {851 id: res.value[0],852 data: res.value[1]853 }854 };855 } else chunk = { data: res.value };856 observer.next({ result: (0, import_objectSpread2.default)((0, import_objectSpread2.default)({}, chunk), {}, { data: transformChunk(chunk.data) }) });857 }858 } catch (_) {859 _usingCtx4.e = _;860 } finally {861 await _usingCtx4.d();862 }863 break;864 } catch (_) {865 _usingCtx3.e = _;866 } finally {867 _usingCtx3.d();868 }869 }870 }).catch((cause) => {871 onErrorCallback(cause);872 observer.error(coerceToTRPCClientError(cause));873 });874 return () => {875 ac.abort();876 };877 });878}879/**880* @deprecated Renamed to `unstable_localLink`. This alias will be removed in a future major release.881*/882const experimental_localLink = unstable_localLink;883 884//#endregion885export { TRPCClientError, TRPCUntypedClient, clientCallTypeToProcedureType, createTRPCClient, createTRPCClientProxy, createTRPCClient as createTRPCProxyClient, createTRPCUntypedClient, createWSClient, experimental_localLink, getFetch, getUntypedClient, httpBatchLink, httpBatchStreamLink, httpLink, httpSubscriptionLink, isFormData, isNonJsonSerializable, isOctetType, isTRPCClientError, jsonEncoder, loggerLink, retryLink, splitLink, unstable_httpBatchStreamLink, unstable_httpSubscriptionLink, unstable_localLink, wsLink };886//# sourceMappingURL=index.mjs.map