CoolFace
Apppublic

AK-21/Graphite-Industrial-Intelligence

sourceHugging Faceupdated 3mo agoView on Hugging Face
0likes
index.mjs886 linesDownload Raw Back to dist
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