EdgeAIG's picture
download
raw
15.6 kB
/**
* Coordinates shared access inside transactions with read and write locks.
*
* A `TxReentrantLock` lets many fibers hold read locks at the same time, or one
* fiber hold a write lock for exclusive access. Lock ownership is tracked by
* fiber, so a fiber that already holds the lock can acquire it again and later
* release each acquisition. Attempts that cannot proceed retry transactionally
* until the lock becomes available. This module includes manual, scoped, and
* wrapper-style operations for read and write locking.
*
* @since 4.0.0
*/
import * as Effect from "./Effect.js";
import * as HashMap from "./HashMap.js";
import { NodeInspectSymbol, toJson } from "./Inspectable.js";
import * as Option from "./Option.js";
import { pipeArguments } from "./Pipeable.js";
import { hasProperty } from "./Predicate.js";
import * as TxRef from "./TxRef.js";
const TypeId = "~effect/transactions/TxReentrantLock";
const emptyState = {
readers: /*#__PURE__*/HashMap.empty(),
writer: /*#__PURE__*/Option.none()
};
const TxReentrantLockProto = {
[NodeInspectSymbol]() {
return toJson(this);
},
toJSON() {
return {
_id: "TxReentrantLock"
};
},
toString() {
return "TxReentrantLock";
},
pipe() {
return pipeArguments(this, arguments);
}
};
// =============================================================================
// Constructors
// =============================================================================
/**
* Creates a new TxReentrantLock.
*
* **Example** (Creating a reentrant lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const isLocked = yield* TxReentrantLock.locked(lock)
* console.log(isLocked) // false
* })
* ```
*
* @category constructors
* @since 2.0.0
*/
export const make = () => Effect.gen(function* () {
const stateRef = yield* TxRef.make(emptyState);
const self = Object.create(TxReentrantLockProto);
self[TypeId] = TypeId;
self.stateRef = stateRef;
return self;
}).pipe(Effect.tx);
// =============================================================================
// Mutations
// =============================================================================
/**
* Acquires a read lock. Blocks if another fiber holds the write lock.
* If the current fiber already holds the write lock, the read lock is granted (reentrancy).
* Returns the current number of read locks held by this fiber.
*
* **Example** (Acquiring a read lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const count = yield* TxReentrantLock.acquireRead(lock)
* console.log(count) // 1
* yield* TxReentrantLock.releaseRead(lock)
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const acquireRead = self => Effect.withFiber(fiber => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
const fiberId = fiber.id;
// If another fiber holds the write lock, retry
if (Option.isSome(state.writer) && state.writer.value[0] !== fiberId) {
return yield* Effect.txRetry;
}
// Grant read lock
const currentCount = Option.getOrElse(HashMap.get(state.readers, fiberId), () => 0);
const newCount = currentCount + 1;
yield* TxRef.set(self.stateRef, {
...state,
readers: HashMap.set(state.readers, fiberId, newCount)
});
return newCount;
}).pipe(Effect.tx));
/**
* Acquires the write lock for the current fiber.
*
* **When to use**
*
* Use to enter an exclusive section manually when `withWriteLock` is not the
* right shape.
*
* **Details**
*
* Blocks if any other fiber holds a read or write lock. If the current fiber
* already holds the write lock, the count is incremented. If the current fiber
* holds a read lock, the write lock is granted as an upgrade.
*
* Returns the current number of write locks held by this fiber.
*
* **Example** (Acquiring a write lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const count = yield* TxReentrantLock.acquireWrite(lock)
* console.log(count) // 1
* yield* TxReentrantLock.releaseWrite(lock)
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const acquireWrite = self => Effect.withFiber(fiber => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
const fiberId = fiber.id;
// If another fiber holds the write lock, retry
if (Option.isSome(state.writer) && state.writer.value[0] !== fiberId) {
return yield* Effect.txRetry;
}
// If other fibers hold read locks, retry
for (const [readerId] of state.readers) {
if (readerId !== fiberId && Option.getOrElse(HashMap.get(state.readers, readerId), () => 0) > 0) {
return yield* Effect.txRetry;
}
}
// Grant write lock
if (Option.isSome(state.writer)) {
// Reentrant: increment write count
const newCount = state.writer.value[1] + 1;
yield* TxRef.set(self.stateRef, {
...state,
writer: Option.some([fiberId, newCount])
});
return newCount;
}
// First write lock acquisition
yield* TxRef.set(self.stateRef, {
...state,
writer: Option.some([fiberId, 1])
});
return 1;
}).pipe(Effect.tx));
/**
* Releases one read lock held by the current fiber.
*
* **When to use**
*
* Use to leave a manually acquired read lock.
*
* **Details**
*
* Returns the remaining number of read locks held by this fiber.
*
* **Example** (Releasing a read lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* yield* TxReentrantLock.acquireRead(lock)
* const remaining = yield* TxReentrantLock.releaseRead(lock)
* console.log(remaining) // 0
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const releaseRead = self => Effect.withFiber(fiber => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
const fiberId = fiber.id;
const currentCount = Option.getOrElse(HashMap.get(state.readers, fiberId), () => 0);
if (currentCount <= 0) return 0;
const newCount = currentCount - 1;
const newReaders = newCount === 0 ? HashMap.remove(state.readers, fiberId) : HashMap.set(state.readers, fiberId, newCount);
yield* TxRef.set(self.stateRef, {
...state,
readers: newReaders
});
return newCount;
}).pipe(Effect.tx));
/**
* Releases one write lock held by the current fiber.
*
* **When to use**
*
* Use to leave a manually acquired write lock.
*
* **Details**
*
* Returns the remaining number of write locks held by this fiber.
*
* **Example** (Releasing a write lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* yield* TxReentrantLock.acquireWrite(lock)
* const remaining = yield* TxReentrantLock.releaseWrite(lock)
* console.log(remaining) // 0
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const releaseWrite = self => Effect.withFiber(fiber => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
const fiberId = fiber.id;
if (Option.isNone(state.writer) || state.writer.value[0] !== fiberId) return 0;
const newCount = state.writer.value[1] - 1;
const newWriter = newCount <= 0 ? Option.none() : Option.some([fiberId, newCount]);
yield* TxRef.set(self.stateRef, {
...state,
writer: newWriter
});
return newCount;
}).pipe(Effect.tx));
/**
* Acquires a read lock for the duration of the scope.
* The lock is automatically released when the scope closes.
*
* **Example** (Holding a scoped read lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
*
* yield* Effect.scoped(
* Effect.gen(function*() {
* yield* TxReentrantLock.readLock(lock)
* // read lock is held for the duration of the scope
* })
* )
* // read lock is released
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const readLock = self => Effect.acquireRelease(acquireRead(self), () => releaseRead(self));
/**
* Acquires a write lock for the duration of the scope.
* The lock is automatically released when the scope closes.
*
* **Example** (Holding a scoped write lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
*
* yield* Effect.scoped(
* Effect.gen(function*() {
* yield* TxReentrantLock.writeLock(lock)
* // write lock is held for the duration of the scope
* })
* )
* // write lock is released
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const writeLock = self => Effect.acquireRelease(acquireWrite(self), () => releaseWrite(self));
/**
* Runs the provided effect while holding a read lock. The lock is automatically
* released after the effect completes, fails, or is interrupted.
*
* **Example** (Running an effect with a read lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const result = yield* TxReentrantLock.withReadLock(
* lock,
* Effect.succeed("read data")
* )
* console.log(result) // "read data"
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const withReadLock = (...args) => {
if (args.length === 1) {
const [effect] = args;
return self => Effect.acquireUseRelease(acquireRead(self), () => effect, () => releaseRead(self));
}
const [self, effect] = args;
return Effect.acquireUseRelease(acquireRead(self), () => effect, () => releaseRead(self));
};
/**
* Runs the provided effect while holding a write lock. The lock is automatically
* released after the effect completes, fails, or is interrupted.
*
* **Example** (Running an effect with a write lock)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const result = yield* TxReentrantLock.withWriteLock(
* lock,
* Effect.succeed("wrote data")
* )
* console.log(result) // "wrote data"
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const withWriteLock = (...args) => {
if (args.length === 1) {
const [effect] = args;
return self => Effect.acquireUseRelease(acquireWrite(self), () => effect, () => releaseWrite(self));
}
const [self, effect] = args;
return Effect.acquireUseRelease(acquireWrite(self), () => effect, () => releaseWrite(self));
};
/**
* Runs an effect while holding a write lock.
*
* **When to use**
*
* Use when you need to run an effect with exclusive write access through a
* `TxReentrantLock` and prefer the concise lock helper.
*
* **Example** (Running an effect with exclusive access)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const result = yield* TxReentrantLock.withLock(
* lock,
* Effect.succeed("exclusive operation")
* )
* console.log(result) // "exclusive operation"
* })
* ```
*
* @category mutations
* @since 2.0.0
*/
export const withLock = withWriteLock;
// =============================================================================
// Getters
// =============================================================================
/**
* Returns the total number of read locks held across all fibers.
*
* **Example** (Counting read locks)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* yield* TxReentrantLock.acquireRead(lock)
* const count = yield* TxReentrantLock.readLocks(lock)
* console.log(count) // 1
* yield* TxReentrantLock.releaseRead(lock)
* })
* ```
*
* @category getters
* @since 2.0.0
*/
export const readLocks = self => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
let total = 0;
for (const [, count] of state.readers) {
total += count;
}
return total;
});
/**
* Returns the number of write locks held (0 or the reentrant count).
*
* **Example** (Counting write locks)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const count = yield* TxReentrantLock.writeLocks(lock)
* console.log(count) // 0
* })
* ```
*
* @category getters
* @since 2.0.0
*/
export const writeLocks = self => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
return Option.isSome(state.writer) ? state.writer.value[1] : 0;
});
/**
* Checks whether the lock is held by any fiber (read or write).
*
* **Example** (Checking whether a lock is held)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const isLocked = yield* TxReentrantLock.locked(lock)
* console.log(isLocked) // false
* })
* ```
*
* @category getters
* @since 2.0.0
*/
export const locked = self => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
return HashMap.size(state.readers) > 0 || Option.isSome(state.writer);
});
/**
* Checks whether any fiber holds a read lock.
*
* **Example** (Checking whether a read lock is held)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const isReadLocked = yield* TxReentrantLock.readLocked(lock)
* console.log(isReadLocked) // false
* })
* ```
*
* @category getters
* @since 2.0.0
*/
export const readLocked = self => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
return HashMap.size(state.readers) > 0;
});
/**
* Checks whether any fiber holds a write lock.
*
* **Example** (Checking whether a write lock is held)
*
* ```ts
* import { Effect, TxReentrantLock } from "effect"
*
* const program = Effect.gen(function*() {
* const lock = yield* TxReentrantLock.make()
* const isWriteLocked = yield* TxReentrantLock.writeLocked(lock)
* console.log(isWriteLocked) // false
* })
* ```
*
* @category getters
* @since 2.0.0
*/
export const writeLocked = self => Effect.gen(function* () {
const state = yield* TxRef.get(self.stateRef);
return Option.isSome(state.writer);
});
// =============================================================================
// Guards
// =============================================================================
/**
* Checks whether the given value is a TxReentrantLock.
*
* **Example** (Checking for TxReentrantLock values)
*
* ```ts
* import { TxReentrantLock } from "effect"
*
* declare const someValue: unknown
*
* if (TxReentrantLock.isTxReentrantLock(someValue)) {
* console.log("This is a TxReentrantLock")
* }
* ```
*
* @category guards
* @since 4.0.0
*/
export const isTxReentrantLock = u => hasProperty(u, TypeId);
//# sourceMappingURL=TxReentrantLock.js.map

Xet Storage Details

Size:
15.6 kB
·
Xet hash:
3bd46374975bf188160a8c2cd5e4cdf0c74de59481d46b8799ee44bb3403447b

Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.