CoolFace
Apppublic

AK-21/Graphite-Industrial-Intelligence

sourceHugging Faceupdated 3mo agoView on Hugging Face
0likes
session.js229 linesDownload Raw Back to node-postgres
1import pg from "pg";2import { NoopCache } from "../cache/core/index.js";3import { entityKind } from "../entity.js";4import { NoopLogger } from "../logger.js";5import { PgTransaction } from "../pg-core/index.js";6import { PgPreparedQuery, PgSession } from "../pg-core/session.js";7import { fillPlaceholders, sql } from "../sql/sql.js";8import { tracer } from "../tracing.js";9import { mapResultRow } from "../utils.js";10const { Pool, types } = pg;11class NodePgPreparedQuery extends PgPreparedQuery {12  constructor(client, queryString, params, logger, cache, queryMetadata, cacheConfig, fields, name, _isResponseInArrayMode, customResultMapper) {13    super({ sql: queryString, params }, cache, queryMetadata, cacheConfig);14    this.client = client;15    this.queryString = queryString;16    this.params = params;17    this.logger = logger;18    this.fields = fields;19    this._isResponseInArrayMode = _isResponseInArrayMode;20    this.customResultMapper = customResultMapper;21    this.rawQueryConfig = {22      name,23      text: queryString,24      types: {25        // @ts-ignore26        getTypeParser: (typeId, format) => {27          if (typeId === types.builtins.TIMESTAMPTZ) {28            return (val) => val;29          }30          if (typeId === types.builtins.TIMESTAMP) {31            return (val) => val;32          }33          if (typeId === types.builtins.DATE) {34            return (val) => val;35          }36          if (typeId === types.builtins.INTERVAL) {37            return (val) => val;38          }39          if (typeId === 1231) {40            return (val) => val;41          }42          if (typeId === 1115) {43            return (val) => val;44          }45          if (typeId === 1185) {46            return (val) => val;47          }48          if (typeId === 1187) {49            return (val) => val;50          }51          if (typeId === 1182) {52            return (val) => val;53          }54          return types.getTypeParser(typeId, format);55        }56      }57    };58    this.queryConfig = {59      name,60      text: queryString,61      rowMode: "array",62      types: {63        // @ts-ignore64        getTypeParser: (typeId, format) => {65          if (typeId === types.builtins.TIMESTAMPTZ) {66            return (val) => val;67          }68          if (typeId === types.builtins.TIMESTAMP) {69            return (val) => val;70          }71          if (typeId === types.builtins.DATE) {72            return (val) => val;73          }74          if (typeId === types.builtins.INTERVAL) {75            return (val) => val;76          }77          if (typeId === 1231) {78            return (val) => val;79          }80          if (typeId === 1115) {81            return (val) => val;82          }83          if (typeId === 1185) {84            return (val) => val;85          }86          if (typeId === 1187) {87            return (val) => val;88          }89          if (typeId === 1182) {90            return (val) => val;91          }92          return types.getTypeParser(typeId, format);93        }94      }95    };96  }97  static [entityKind] = "NodePgPreparedQuery";98  rawQueryConfig;99  queryConfig;100  async execute(placeholderValues = {}) {101    return tracer.startActiveSpan("drizzle.execute", async () => {102      const params = fillPlaceholders(this.params, placeholderValues);103      this.logger.logQuery(this.rawQueryConfig.text, params);104      const { fields, rawQueryConfig: rawQuery, client, queryConfig: query, joinsNotNullableMap, customResultMapper } = this;105      if (!fields && !customResultMapper) {106        return tracer.startActiveSpan("drizzle.driver.execute", async (span) => {107          span?.setAttributes({108            "drizzle.query.name": rawQuery.name,109            "drizzle.query.text": rawQuery.text,110            "drizzle.query.params": JSON.stringify(params)111          });112          return this.queryWithCache(rawQuery.text, params, async () => {113            return await client.query(rawQuery, params);114          });115        });116      }117      const result = await tracer.startActiveSpan("drizzle.driver.execute", (span) => {118        span?.setAttributes({119          "drizzle.query.name": query.name,120          "drizzle.query.text": query.text,121          "drizzle.query.params": JSON.stringify(params)122        });123        return this.queryWithCache(query.text, params, async () => {124          return await client.query(query, params);125        });126      });127      return tracer.startActiveSpan("drizzle.mapResponse", () => {128        return customResultMapper ? customResultMapper(result.rows) : result.rows.map((row) => mapResultRow(fields, row, joinsNotNullableMap));129      });130    });131  }132  all(placeholderValues = {}) {133    return tracer.startActiveSpan("drizzle.execute", () => {134      const params = fillPlaceholders(this.params, placeholderValues);135      this.logger.logQuery(this.rawQueryConfig.text, params);136      return tracer.startActiveSpan("drizzle.driver.execute", (span) => {137        span?.setAttributes({138          "drizzle.query.name": this.rawQueryConfig.name,139          "drizzle.query.text": this.rawQueryConfig.text,140          "drizzle.query.params": JSON.stringify(params)141        });142        return this.queryWithCache(this.rawQueryConfig.text, params, async () => {143          return this.client.query(this.rawQueryConfig, params);144        }).then((result) => result.rows);145      });146    });147  }148  /** @internal */149  isResponseInArrayMode() {150    return this._isResponseInArrayMode;151  }152}153class NodePgSession extends PgSession {154  constructor(client, dialect, schema, options = {}) {155    super(dialect);156    this.client = client;157    this.schema = schema;158    this.options = options;159    this.logger = options.logger ?? new NoopLogger();160    this.cache = options.cache ?? new NoopCache();161  }162  static [entityKind] = "NodePgSession";163  logger;164  cache;165  prepareQuery(query, fields, name, isResponseInArrayMode, customResultMapper, queryMetadata, cacheConfig) {166    return new NodePgPreparedQuery(167      this.client,168      query.sql,169      query.params,170      this.logger,171      this.cache,172      queryMetadata,173      cacheConfig,174      fields,175      name,176      isResponseInArrayMode,177      customResultMapper178    );179  }180  async transaction(transaction, config) {181    const isPool = this.client instanceof Pool || Object.getPrototypeOf(this.client).constructor.name.includes("Pool");182    const session = isPool ? new NodePgSession(await this.client.connect(), this.dialect, this.schema, this.options) : this;183    const tx = new NodePgTransaction(this.dialect, session, this.schema);184    await tx.execute(sql`begin${config ? sql` ${tx.getTransactionConfigSQL(config)}` : void 0}`);185    try {186      const result = await transaction(tx);187      await tx.execute(sql`commit`);188      return result;189    } catch (error) {190      await tx.execute(sql`rollback`);191      throw error;192    } finally {193      if (isPool) session.client.release();194    }195  }196  async count(sql2) {197    const res = await this.execute(sql2);198    return Number(199      res["rows"][0]["count"]200    );201  }202}203class NodePgTransaction extends PgTransaction {204  static [entityKind] = "NodePgTransaction";205  async transaction(transaction) {206    const savepointName = `sp${this.nestedIndex + 1}`;207    const tx = new NodePgTransaction(208      this.dialect,209      this.session,210      this.schema,211      this.nestedIndex + 1212    );213    await tx.execute(sql.raw(`savepoint ${savepointName}`));214    try {215      const result = await transaction(tx);216      await tx.execute(sql.raw(`release savepoint ${savepointName}`));217      return result;218    } catch (err) {219      await tx.execute(sql.raw(`rollback to savepoint ${savepointName}`));220      throw err;221    }222  }223}224export {225  NodePgPreparedQuery,226  NodePgSession,227  NodePgTransaction228};229//# sourceMappingURL=session.js.map