mirror of
https://github.com/Start9Labs/start-os.git
synced 2026-03-30 04:01:58 +00:00
refactor: consolidate SDK Watchable with generic map/eq and rename call to fetch
This commit is contained in:
@@ -12,7 +12,7 @@ export class GetContainerIp extends Watchable<string> {
|
|||||||
super(effects)
|
super(effects)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getContainerIp({ ...this.opts, callback })
|
return this.effects.getContainerIp({ ...this.opts, callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ export class GetHostInfo extends Watchable<Host | null> {
|
|||||||
super(effects)
|
super(effects)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getHostInfo({ ...this.opts, callback })
|
return this.effects.getHostInfo({ ...this.opts, callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -8,7 +8,7 @@ export class GetOutboundGateway extends Watchable<string> {
|
|||||||
super(effects)
|
super(effects)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getOutboundGateway({ callback })
|
return this.effects.getOutboundGateway({ callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,18 +1,47 @@
|
|||||||
import { Effects } from '../Effects'
|
import { Effects } from '../Effects'
|
||||||
import { Manifest, PackageId } from '../osBindings'
|
import { Manifest, PackageId } from '../osBindings'
|
||||||
|
import { deepEqual } from './deepEqual'
|
||||||
import { Watchable } from './Watchable'
|
import { Watchable } from './Watchable'
|
||||||
|
|
||||||
export class GetServiceManifest extends Watchable<Manifest> {
|
export class GetServiceManifest<
|
||||||
|
Mapped = Manifest | null,
|
||||||
|
> extends Watchable<Manifest | null, Mapped> {
|
||||||
protected readonly label = 'GetServiceManifest'
|
protected readonly label = 'GetServiceManifest'
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
effects: Effects,
|
effects: Effects,
|
||||||
readonly opts: { packageId: PackageId },
|
readonly opts: { packageId: PackageId },
|
||||||
|
options?: {
|
||||||
|
map?: (value: Manifest | null) => Mapped
|
||||||
|
eq?: (a: Mapped, b: Mapped) => boolean
|
||||||
|
},
|
||||||
) {
|
) {
|
||||||
super(effects)
|
super(effects, options)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getServiceManifest({ ...this.opts, callback })
|
return this.effects.getServiceManifest({ ...this.opts, callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function getServiceManifest(
|
||||||
|
effects: Effects,
|
||||||
|
packageId: PackageId,
|
||||||
|
): GetServiceManifest<Manifest | null>
|
||||||
|
export function getServiceManifest<Mapped>(
|
||||||
|
effects: Effects,
|
||||||
|
packageId: PackageId,
|
||||||
|
map: (manifest: Manifest | null) => Mapped,
|
||||||
|
eq?: (a: Mapped, b: Mapped) => boolean,
|
||||||
|
): GetServiceManifest<Mapped>
|
||||||
|
export function getServiceManifest<Mapped>(
|
||||||
|
effects: Effects,
|
||||||
|
packageId: PackageId,
|
||||||
|
map?: (manifest: Manifest | null) => Mapped,
|
||||||
|
eq?: (a: Mapped, b: Mapped) => boolean,
|
||||||
|
): GetServiceManifest<Mapped> {
|
||||||
|
return new GetServiceManifest<Mapped>(effects, { packageId }, {
|
||||||
|
map: map ?? ((a) => a as Mapped),
|
||||||
|
eq: eq ?? ((a, b) => deepEqual(a, b)),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ export class GetSslCertificate extends Watchable<[string, string, string]> {
|
|||||||
super(effects)
|
super(effects)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getSslCertificate({ ...this.opts, callback })
|
return this.effects.getSslCertificate({ ...this.opts, callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ export class GetStatus extends Watchable<StatusInfo | null> {
|
|||||||
super(effects)
|
super(effects)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getStatus({ ...this.opts, callback })
|
return this.effects.getStatus({ ...this.opts, callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ export class GetSystemSmtp extends Watchable<T.SmtpValue | null> {
|
|||||||
super(effects)
|
super(effects)
|
||||||
}
|
}
|
||||||
|
|
||||||
protected call(callback?: () => void) {
|
protected fetch(callback?: () => void) {
|
||||||
return this.effects.getSystemSmtp({ callback })
|
return this.effects.getSystemSmtp({ callback })
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,55 +1,124 @@
|
|||||||
import { Effects } from '../Effects'
|
import { Effects } from '../Effects'
|
||||||
import { AbortedError } from './AbortedError'
|
import { AbortedError } from './AbortedError'
|
||||||
|
import { deepEqual } from './deepEqual'
|
||||||
import { DropGenerator, DropPromise } from './Drop'
|
import { DropGenerator, DropPromise } from './Drop'
|
||||||
|
|
||||||
export abstract class Watchable<T> {
|
export abstract class Watchable<Raw, Mapped = Raw> {
|
||||||
constructor(readonly effects: Effects) {}
|
protected readonly mapFn: (value: Raw) => Mapped
|
||||||
|
protected readonly eqFn: (a: Mapped, b: Mapped) => boolean
|
||||||
|
|
||||||
protected abstract call(callback?: () => void): Promise<T>
|
constructor(
|
||||||
|
readonly effects: Effects,
|
||||||
|
options?: {
|
||||||
|
map?: (value: Raw) => Mapped
|
||||||
|
eq?: (a: Mapped, b: Mapped) => boolean
|
||||||
|
},
|
||||||
|
) {
|
||||||
|
this.mapFn = options?.map ?? ((a) => a as unknown as Mapped)
|
||||||
|
this.eqFn = options?.eq ?? ((a, b) => deepEqual(a, b))
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Fetch the current value, optionally registering a callback for change notification.
|
||||||
|
* The callback should be invoked when the underlying data changes.
|
||||||
|
*/
|
||||||
|
protected abstract fetch(callback?: () => void): Promise<Raw>
|
||||||
protected abstract readonly label: string
|
protected abstract readonly label: string
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Returns the value. Reruns the context from which it has been called if the underlying value changes
|
* Produce a stream of raw values. Default implementation uses fetch() with
|
||||||
|
* effects callback in a loop. Override for custom subscription mechanisms
|
||||||
|
* (e.g. fs.watch).
|
||||||
*/
|
*/
|
||||||
const(): Promise<T> {
|
protected async *produce(abort: AbortSignal): AsyncGenerator<Raw, void> {
|
||||||
return this.call(
|
|
||||||
this.effects.constRetry &&
|
|
||||||
(() => this.effects.constRetry && this.effects.constRetry()),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Returns the value. Does nothing if the value changes
|
|
||||||
*/
|
|
||||||
once(): Promise<T> {
|
|
||||||
return this.call()
|
|
||||||
}
|
|
||||||
|
|
||||||
private async *watchGen(abort?: AbortSignal) {
|
|
||||||
const resolveCell = { resolve: () => {} }
|
const resolveCell = { resolve: () => {} }
|
||||||
this.effects.onLeaveContext(() => {
|
this.effects.onLeaveContext(() => {
|
||||||
resolveCell.resolve()
|
resolveCell.resolve()
|
||||||
})
|
})
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
abort.addEventListener('abort', () => resolveCell.resolve())
|
||||||
while (this.effects.isInContext && !abort?.aborted) {
|
while (this.effects.isInContext && !abort.aborted) {
|
||||||
let callback: () => void = () => {}
|
let callback: () => void = () => {}
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
const waitForNext = new Promise<void>((resolve) => {
|
||||||
callback = resolve
|
callback = resolve
|
||||||
resolveCell.resolve = resolve
|
resolveCell.resolve = resolve
|
||||||
})
|
})
|
||||||
yield await this.call(() => callback())
|
yield await this.fetch(() => callback())
|
||||||
await waitForNext
|
await waitForNext
|
||||||
}
|
}
|
||||||
return new Promise<never>((_, rej) => rej(new AbortedError()))
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Lifecycle hook called when const() registers a subscription.
|
||||||
|
* Return a cleanup function to be called when the subscription ends.
|
||||||
|
* Override for side effects like FileHelper's consts tracking.
|
||||||
|
*/
|
||||||
|
protected onConstRegistered(_value: Mapped): (() => void) | void {}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Internal generator that maps raw values and deduplicates using eq.
|
||||||
|
*/
|
||||||
|
private async *watchGen(
|
||||||
|
abort: AbortSignal,
|
||||||
|
): AsyncGenerator<Mapped, void, unknown> {
|
||||||
|
let prev: { value: Mapped } | null = null
|
||||||
|
for await (const raw of this.produce(abort)) {
|
||||||
|
if (abort.aborted) return
|
||||||
|
const mapped = this.mapFn(raw)
|
||||||
|
if (!prev || !this.eqFn(prev.value, mapped)) {
|
||||||
|
prev = { value: mapped }
|
||||||
|
yield mapped
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Returns the value. Reruns the context from which it has been called if the underlying value changes
|
||||||
|
*/
|
||||||
|
async const(): Promise<Mapped> {
|
||||||
|
const abort = new AbortController()
|
||||||
|
const gen = this.watchGen(abort.signal)
|
||||||
|
const res = await gen.next()
|
||||||
|
const value = res.value as Mapped
|
||||||
|
if (this.effects.constRetry) {
|
||||||
|
const constRetry = this.effects.constRetry
|
||||||
|
const cleanup = this.onConstRegistered(value)
|
||||||
|
gen.next().then(
|
||||||
|
() => {
|
||||||
|
abort.abort()
|
||||||
|
cleanup?.()
|
||||||
|
constRetry()
|
||||||
|
},
|
||||||
|
() => {
|
||||||
|
abort.abort()
|
||||||
|
cleanup?.()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
abort.abort()
|
||||||
|
}
|
||||||
|
return value
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Returns the value. Does nothing if the value changes
|
||||||
|
*/
|
||||||
|
async once(): Promise<Mapped> {
|
||||||
|
return this.mapFn(await this.fetch())
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Watches the value. Returns an async iterator that yields whenever the value changes
|
* Watches the value. Returns an async iterator that yields whenever the value changes
|
||||||
*/
|
*/
|
||||||
watch(abort?: AbortSignal): AsyncGenerator<T, never, unknown> {
|
watch(abort?: AbortSignal): AsyncGenerator<Mapped, never, unknown> {
|
||||||
const ctrl = new AbortController()
|
const ctrl = new AbortController()
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
abort?.addEventListener('abort', () => ctrl.abort())
|
||||||
return DropGenerator.of(this.watchGen(ctrl.signal), () => ctrl.abort())
|
return DropGenerator.of(
|
||||||
|
(async function* (gen): AsyncGenerator<Mapped, never, unknown> {
|
||||||
|
yield* gen
|
||||||
|
throw new AbortedError()
|
||||||
|
})(this.watchGen(ctrl.signal)),
|
||||||
|
() => ctrl.abort(),
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -57,13 +126,13 @@ export abstract class Watchable<T> {
|
|||||||
*/
|
*/
|
||||||
onChange(
|
onChange(
|
||||||
callback: (
|
callback: (
|
||||||
value: T | undefined,
|
value: Mapped | undefined,
|
||||||
error?: Error,
|
error?: Error,
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
||||||
) {
|
) {
|
||||||
;(async () => {
|
;(async () => {
|
||||||
const ctrl = new AbortController()
|
const ctrl = new AbortController()
|
||||||
for await (const value of this.watch(ctrl.signal)) {
|
for await (const value of this.watchGen(ctrl.signal)) {
|
||||||
try {
|
try {
|
||||||
const res = await callback(value)
|
const res = await callback(value)
|
||||||
if (res.cancel) {
|
if (res.cancel) {
|
||||||
@@ -90,7 +159,7 @@ export abstract class Watchable<T> {
|
|||||||
/**
|
/**
|
||||||
* Watches the value. Returns when the predicate is true
|
* Watches the value. Returns when the predicate is true
|
||||||
*/
|
*/
|
||||||
waitFor(pred: (value: T) => boolean): Promise<T> {
|
waitFor(pred: (value: Mapped) => boolean): Promise<Mapped> {
|
||||||
const ctrl = new AbortController()
|
const ctrl = new AbortController()
|
||||||
return DropPromise.of(
|
return DropPromise.of(
|
||||||
Promise.resolve().then(async () => {
|
Promise.resolve().then(async () => {
|
||||||
|
|||||||
@@ -8,11 +8,10 @@ import {
|
|||||||
HostnameInfo,
|
HostnameInfo,
|
||||||
} from '../types'
|
} from '../types'
|
||||||
import { Effects } from '../Effects'
|
import { Effects } from '../Effects'
|
||||||
import { AbortedError } from './AbortedError'
|
|
||||||
import { DropGenerator, DropPromise } from './Drop'
|
|
||||||
import { IpAddress, IPV6_LINK_LOCAL } from './ip'
|
import { IpAddress, IPV6_LINK_LOCAL } from './ip'
|
||||||
import { deepEqual } from './deepEqual'
|
import { deepEqual } from './deepEqual'
|
||||||
import { once } from './once'
|
import { once } from './once'
|
||||||
|
import { Watchable } from './Watchable'
|
||||||
|
|
||||||
export type UrlString = string
|
export type UrlString = string
|
||||||
export type HostId = string
|
export type HostId = string
|
||||||
@@ -440,136 +439,29 @@ const makeInterfaceFilled = async ({
|
|||||||
return interfaceFilled
|
return interfaceFilled
|
||||||
}
|
}
|
||||||
|
|
||||||
export class GetServiceInterface<Mapped = ServiceInterfaceFilled | null> {
|
export class GetServiceInterface<
|
||||||
|
Mapped = ServiceInterfaceFilled | null,
|
||||||
|
> extends Watchable<ServiceInterfaceFilled | null, Mapped> {
|
||||||
|
protected readonly label = 'GetServiceInterface'
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
readonly effects: Effects,
|
effects: Effects,
|
||||||
readonly opts: { id: string; packageId?: string },
|
readonly opts: { id: string; packageId?: string },
|
||||||
readonly map: (interfaces: ServiceInterfaceFilled | null) => Mapped,
|
options?: {
|
||||||
readonly eq: (a: Mapped, b: Mapped) => boolean,
|
map?: (value: ServiceInterfaceFilled | null) => Mapped
|
||||||
) {}
|
eq?: (a: Mapped, b: Mapped) => boolean
|
||||||
|
},
|
||||||
/**
|
|
||||||
* Returns the requested service interface. Reruns the context from which it has been called if the underlying value changes
|
|
||||||
*/
|
|
||||||
async const() {
|
|
||||||
let abort = new AbortController()
|
|
||||||
const watch = this.watch(abort.signal)
|
|
||||||
const res = await watch.next()
|
|
||||||
if (this.effects.constRetry) {
|
|
||||||
watch
|
|
||||||
.next()
|
|
||||||
.then(() => {
|
|
||||||
abort.abort()
|
|
||||||
this.effects.constRetry && this.effects.constRetry()
|
|
||||||
})
|
|
||||||
.catch()
|
|
||||||
}
|
|
||||||
return res.value
|
|
||||||
}
|
|
||||||
/**
|
|
||||||
* Returns the requested service interface. Does nothing if the value changes
|
|
||||||
*/
|
|
||||||
async once() {
|
|
||||||
const { id, packageId } = this.opts
|
|
||||||
const interfaceFilled = await makeInterfaceFilled({
|
|
||||||
effects: this.effects,
|
|
||||||
id,
|
|
||||||
packageId,
|
|
||||||
})
|
|
||||||
|
|
||||||
return this.map(interfaceFilled)
|
|
||||||
}
|
|
||||||
|
|
||||||
private async *watchGen(abort?: AbortSignal) {
|
|
||||||
let prev = null as { value: Mapped } | null
|
|
||||||
const { id, packageId } = this.opts
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
this.effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
|
||||||
while (this.effects.isInContext && !abort?.aborted) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
const next = this.map(
|
|
||||||
await makeInterfaceFilled({
|
|
||||||
effects: this.effects,
|
|
||||||
id,
|
|
||||||
packageId,
|
|
||||||
callback,
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
if (!prev || !this.eq(prev.value, next)) {
|
|
||||||
yield next
|
|
||||||
}
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
return new Promise<never>((_, rej) => rej(new AbortedError()))
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the requested service interface. Returns an async iterator that yields whenever the value changes
|
|
||||||
*/
|
|
||||||
watch(abort?: AbortSignal): AsyncGenerator<Mapped, never, unknown> {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
return DropGenerator.of(this.watchGen(ctrl.signal), () => ctrl.abort())
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the requested service interface. Takes a custom callback function to run whenever the value changes
|
|
||||||
*/
|
|
||||||
onChange(
|
|
||||||
callback: (
|
|
||||||
value: Mapped | null,
|
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
|
||||||
) {
|
) {
|
||||||
;(async () => {
|
super(effects, options)
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of this.watch(ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) {
|
|
||||||
ctrl.abort()
|
|
||||||
break
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetServiceInterface.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetServiceInterface.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
protected fetch(callback?: () => void) {
|
||||||
* Watches the requested service interface. Returns when the predicate is true
|
return makeInterfaceFilled({
|
||||||
*/
|
effects: this.effects,
|
||||||
waitFor(pred: (value: Mapped) => boolean): Promise<Mapped> {
|
id: this.opts.id,
|
||||||
const ctrl = new AbortController()
|
packageId: this.opts.packageId,
|
||||||
return DropPromise.of(
|
callback,
|
||||||
Promise.resolve().then(async () => {
|
})
|
||||||
for await (const next of this.watchGen(ctrl.signal)) {
|
|
||||||
if (pred(next)) {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
throw new Error('context left before predicate passed')
|
|
||||||
}),
|
|
||||||
() => ctrl.abort(),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -589,12 +481,10 @@ export function getOwnServiceInterface<Mapped>(
|
|||||||
map?: (interfaces: ServiceInterfaceFilled | null) => Mapped,
|
map?: (interfaces: ServiceInterfaceFilled | null) => Mapped,
|
||||||
eq?: (a: Mapped, b: Mapped) => boolean,
|
eq?: (a: Mapped, b: Mapped) => boolean,
|
||||||
): GetServiceInterface<Mapped> {
|
): GetServiceInterface<Mapped> {
|
||||||
return new GetServiceInterface(
|
return new GetServiceInterface<Mapped>(effects, { id }, {
|
||||||
effects,
|
map: map ?? ((a) => a as Mapped),
|
||||||
{ id },
|
eq: eq ?? ((a, b) => deepEqual(a, b)),
|
||||||
map ?? ((a) => a as Mapped),
|
})
|
||||||
eq ?? ((a, b) => deepEqual(a, b)),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export function getServiceInterface(
|
export function getServiceInterface(
|
||||||
@@ -613,10 +503,8 @@ export function getServiceInterface<Mapped>(
|
|||||||
map?: (interfaces: ServiceInterfaceFilled | null) => Mapped,
|
map?: (interfaces: ServiceInterfaceFilled | null) => Mapped,
|
||||||
eq?: (a: Mapped, b: Mapped) => boolean,
|
eq?: (a: Mapped, b: Mapped) => boolean,
|
||||||
): GetServiceInterface<Mapped> {
|
): GetServiceInterface<Mapped> {
|
||||||
return new GetServiceInterface(
|
return new GetServiceInterface<Mapped>(effects, opts, {
|
||||||
effects,
|
map: map ?? ((a) => a as Mapped),
|
||||||
opts,
|
eq: eq ?? ((a, b) => deepEqual(a, b)),
|
||||||
map ?? ((a) => a as Mapped),
|
})
|
||||||
eq ?? ((a, b) => deepEqual(a, b)),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,9 +1,8 @@
|
|||||||
import { Effects } from '../Effects'
|
import { Effects } from '../Effects'
|
||||||
import { PackageId } from '../osBindings'
|
import { PackageId } from '../osBindings'
|
||||||
import { AbortedError } from './AbortedError'
|
|
||||||
import { deepEqual } from './deepEqual'
|
import { deepEqual } from './deepEqual'
|
||||||
import { DropGenerator, DropPromise } from './Drop'
|
|
||||||
import { ServiceInterfaceFilled, filledAddress } from './getServiceInterface'
|
import { ServiceInterfaceFilled, filledAddress } from './getServiceInterface'
|
||||||
|
import { Watchable } from './Watchable'
|
||||||
|
|
||||||
const makeManyInterfaceFilled = async ({
|
const makeManyInterfaceFilled = async ({
|
||||||
effects,
|
effects,
|
||||||
@@ -40,139 +39,34 @@ const makeManyInterfaceFilled = async ({
|
|||||||
return serviceInterfacesFilled
|
return serviceInterfacesFilled
|
||||||
}
|
}
|
||||||
|
|
||||||
export class GetServiceInterfaces<Mapped = ServiceInterfaceFilled[]> {
|
export class GetServiceInterfaces<
|
||||||
|
Mapped = ServiceInterfaceFilled[],
|
||||||
|
> extends Watchable<ServiceInterfaceFilled[], Mapped> {
|
||||||
|
protected readonly label = 'GetServiceInterfaces'
|
||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
readonly effects: Effects,
|
effects: Effects,
|
||||||
readonly opts: { packageId?: string },
|
readonly opts: { packageId?: string },
|
||||||
readonly map: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
options?: {
|
||||||
readonly eq: (a: Mapped, b: Mapped) => boolean,
|
map?: (value: ServiceInterfaceFilled[]) => Mapped
|
||||||
) {}
|
eq?: (a: Mapped, b: Mapped) => boolean
|
||||||
|
},
|
||||||
/**
|
|
||||||
* Returns the service interfaces for the package. Reruns the context from which it has been called if the underlying value changes
|
|
||||||
*/
|
|
||||||
async const() {
|
|
||||||
let abort = new AbortController()
|
|
||||||
const watch = this.watch(abort.signal)
|
|
||||||
const res = await watch.next()
|
|
||||||
if (this.effects.constRetry) {
|
|
||||||
watch
|
|
||||||
.next()
|
|
||||||
.then(() => {
|
|
||||||
abort.abort()
|
|
||||||
this.effects.constRetry && this.effects.constRetry()
|
|
||||||
})
|
|
||||||
.catch()
|
|
||||||
}
|
|
||||||
return res.value
|
|
||||||
}
|
|
||||||
/**
|
|
||||||
* Returns the service interfaces for the package. Does nothing if the value changes
|
|
||||||
*/
|
|
||||||
async once() {
|
|
||||||
const { packageId } = this.opts
|
|
||||||
const interfaceFilled: ServiceInterfaceFilled[] =
|
|
||||||
await makeManyInterfaceFilled({
|
|
||||||
effects: this.effects,
|
|
||||||
packageId,
|
|
||||||
})
|
|
||||||
|
|
||||||
return this.map(interfaceFilled)
|
|
||||||
}
|
|
||||||
|
|
||||||
private async *watchGen(abort?: AbortSignal) {
|
|
||||||
let prev = null as { value: Mapped } | null
|
|
||||||
const { packageId } = this.opts
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
this.effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
|
||||||
while (this.effects.isInContext && !abort?.aborted) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
const next = this.map(
|
|
||||||
await makeManyInterfaceFilled({
|
|
||||||
effects: this.effects,
|
|
||||||
packageId,
|
|
||||||
callback,
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
if (!prev || !this.eq(prev.value, next)) {
|
|
||||||
yield next
|
|
||||||
}
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
return new Promise<never>((_, rej) => rej(new AbortedError()))
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the service interfaces for the package. Returns an async iterator that yields whenever the value changes
|
|
||||||
*/
|
|
||||||
watch(abort?: AbortSignal): AsyncGenerator<Mapped, never, unknown> {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
return DropGenerator.of(this.watchGen(ctrl.signal), () => ctrl.abort())
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the service interfaces for the package. Takes a custom callback function to run whenever the value changes
|
|
||||||
*/
|
|
||||||
onChange(
|
|
||||||
callback: (
|
|
||||||
value: Mapped | null,
|
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
|
||||||
) {
|
) {
|
||||||
;(async () => {
|
super(effects, options)
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of this.watch(ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) {
|
|
||||||
ctrl.abort()
|
|
||||||
break
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetServiceInterfaces.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetServiceInterfaces.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
protected fetch(callback?: () => void) {
|
||||||
* Watches the service interfaces for the package. Returns when the predicate is true
|
return makeManyInterfaceFilled({
|
||||||
*/
|
effects: this.effects,
|
||||||
waitFor(pred: (value: Mapped) => boolean): Promise<Mapped> {
|
packageId: this.opts.packageId,
|
||||||
const ctrl = new AbortController()
|
callback,
|
||||||
return DropPromise.of(
|
})
|
||||||
Promise.resolve().then(async () => {
|
|
||||||
for await (const next of this.watchGen(ctrl.signal)) {
|
|
||||||
if (pred(next)) {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
throw new Error('context left before predicate passed')
|
|
||||||
}),
|
|
||||||
() => ctrl.abort(),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export function getOwnServiceInterfaces(effects: Effects): GetServiceInterfaces
|
export function getOwnServiceInterfaces(
|
||||||
|
effects: Effects,
|
||||||
|
): GetServiceInterfaces
|
||||||
export function getOwnServiceInterfaces<Mapped>(
|
export function getOwnServiceInterfaces<Mapped>(
|
||||||
effects: Effects,
|
effects: Effects,
|
||||||
map: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
map: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
||||||
@@ -183,12 +77,10 @@ export function getOwnServiceInterfaces<Mapped>(
|
|||||||
map?: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
map?: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
||||||
eq?: (a: Mapped, b: Mapped) => boolean,
|
eq?: (a: Mapped, b: Mapped) => boolean,
|
||||||
): GetServiceInterfaces<Mapped> {
|
): GetServiceInterfaces<Mapped> {
|
||||||
return new GetServiceInterfaces(
|
return new GetServiceInterfaces<Mapped>(effects, {}, {
|
||||||
effects,
|
map: map ?? ((a) => a as Mapped),
|
||||||
{},
|
eq: eq ?? ((a, b) => deepEqual(a, b)),
|
||||||
map ?? ((a) => a as Mapped),
|
})
|
||||||
eq ?? ((a, b) => deepEqual(a, b)),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export function getServiceInterfaces(
|
export function getServiceInterfaces(
|
||||||
@@ -207,10 +99,8 @@ export function getServiceInterfaces<Mapped>(
|
|||||||
map?: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
map?: (interfaces: ServiceInterfaceFilled[]) => Mapped,
|
||||||
eq?: (a: Mapped, b: Mapped) => boolean,
|
eq?: (a: Mapped, b: Mapped) => boolean,
|
||||||
): GetServiceInterfaces<Mapped> {
|
): GetServiceInterfaces<Mapped> {
|
||||||
return new GetServiceInterfaces(
|
return new GetServiceInterfaces<Mapped>(effects, opts, {
|
||||||
effects,
|
map: map ?? ((a) => a as Mapped),
|
||||||
opts,
|
eq: eq ?? ((a, b) => deepEqual(a, b)),
|
||||||
map ?? ((a) => a as Mapped),
|
})
|
||||||
eq ?? ((a, b) => deepEqual(a, b)),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -19,7 +19,7 @@ export { Watchable } from './Watchable'
|
|||||||
export { GetContainerIp } from './GetContainerIp'
|
export { GetContainerIp } from './GetContainerIp'
|
||||||
export { GetHostInfo } from './GetHostInfo'
|
export { GetHostInfo } from './GetHostInfo'
|
||||||
export { GetOutboundGateway } from './GetOutboundGateway'
|
export { GetOutboundGateway } from './GetOutboundGateway'
|
||||||
export { GetServiceManifest } from './GetServiceManifest'
|
export { GetServiceManifest, getServiceManifest } from './GetServiceManifest'
|
||||||
export { GetSslCertificate } from './GetSslCertificate'
|
export { GetSslCertificate } from './GetSslCertificate'
|
||||||
export { GetStatus } from './GetStatus'
|
export { GetStatus } from './GetStatus'
|
||||||
export { GetSystemSmtp } from './GetSystemSmtp'
|
export { GetSystemSmtp } from './GetSystemSmtp'
|
||||||
|
|||||||
@@ -59,7 +59,8 @@ import {
|
|||||||
setupOnInit,
|
setupOnInit,
|
||||||
setupOnUninit,
|
setupOnUninit,
|
||||||
} from '../../base/lib/inits'
|
} from '../../base/lib/inits'
|
||||||
import { DropGenerator } from '../../base/lib/util/Drop'
|
import { GetContainerIp } from '../../base/lib/util/GetContainerIp'
|
||||||
|
import { GetStatus } from '../../base/lib/util/GetStatus'
|
||||||
import {
|
import {
|
||||||
getOwnServiceInterface,
|
getOwnServiceInterface,
|
||||||
ServiceInterfaceFilled,
|
ServiceInterfaceFilled,
|
||||||
@@ -257,90 +258,7 @@ export class StartSdk<Manifest extends T.SDKManifest> {
|
|||||||
Parameters<T.Effects['getContainerIp']>[0],
|
Parameters<T.Effects['getContainerIp']>[0],
|
||||||
'callback'
|
'callback'
|
||||||
> = {},
|
> = {},
|
||||||
) => {
|
) => new GetContainerIp(effects, options),
|
||||||
async function* watch(abort?: AbortSignal) {
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
|
||||||
while (effects.isInContext && !abort?.aborted) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
yield await effects.getContainerIp({ ...options, callback })
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return {
|
|
||||||
const: () =>
|
|
||||||
effects.getContainerIp({
|
|
||||||
...options,
|
|
||||||
callback:
|
|
||||||
effects.constRetry &&
|
|
||||||
(() => effects.constRetry && effects.constRetry()),
|
|
||||||
}),
|
|
||||||
once: () => effects.getContainerIp(options),
|
|
||||||
watch: (abort?: AbortSignal) => {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
return DropGenerator.of(watch(ctrl.signal), () => ctrl.abort())
|
|
||||||
},
|
|
||||||
onChange: (
|
|
||||||
callback: (
|
|
||||||
value: string | null,
|
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
|
||||||
) => {
|
|
||||||
;(async () => {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of watch(ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) {
|
|
||||||
ctrl.abort()
|
|
||||||
break
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ getContainerIp.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ getContainerIp.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
},
|
|
||||||
waitFor: async (pred: (value: string | null) => boolean) => {
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
while (effects.isInContext) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
const res = await effects.getContainerIp({ ...options, callback })
|
|
||||||
if (pred(res)) {
|
|
||||||
resolveCell.resolve()
|
|
||||||
return res
|
|
||||||
}
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
return null
|
|
||||||
},
|
|
||||||
}
|
|
||||||
},
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Get the service's current status with reactive subscription support.
|
* Get the service's current status with reactive subscription support.
|
||||||
@@ -355,90 +273,7 @@ export class StartSdk<Manifest extends T.SDKManifest> {
|
|||||||
getStatus: (
|
getStatus: (
|
||||||
effects: T.Effects,
|
effects: T.Effects,
|
||||||
options: Omit<Parameters<T.Effects['getStatus']>[0], 'callback'> = {},
|
options: Omit<Parameters<T.Effects['getStatus']>[0], 'callback'> = {},
|
||||||
) => {
|
) => new GetStatus(effects, options),
|
||||||
async function* watch(abort?: AbortSignal) {
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
|
||||||
while (effects.isInContext && !abort?.aborted) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
yield await effects.getStatus({ ...options, callback })
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return {
|
|
||||||
const: () =>
|
|
||||||
effects.getStatus({
|
|
||||||
...options,
|
|
||||||
callback:
|
|
||||||
effects.constRetry &&
|
|
||||||
(() => effects.constRetry && effects.constRetry()),
|
|
||||||
}),
|
|
||||||
once: () => effects.getStatus(options),
|
|
||||||
watch: (abort?: AbortSignal) => {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
return DropGenerator.of(watch(ctrl.signal), () => ctrl.abort())
|
|
||||||
},
|
|
||||||
onChange: (
|
|
||||||
callback: (
|
|
||||||
value: T.StatusInfo | null,
|
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
|
||||||
) => {
|
|
||||||
;(async () => {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of watch(ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) {
|
|
||||||
ctrl.abort()
|
|
||||||
break
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ getStatus.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ getStatus.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
},
|
|
||||||
waitFor: async (pred: (value: T.StatusInfo | null) => boolean) => {
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
while (effects.isInContext) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
const res = await effects.getStatus({ ...options, callback })
|
|
||||||
if (pred(res)) {
|
|
||||||
resolveCell.resolve()
|
|
||||||
return res
|
|
||||||
}
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
return null
|
|
||||||
},
|
|
||||||
}
|
|
||||||
},
|
|
||||||
|
|
||||||
MultiHost: {
|
MultiHost: {
|
||||||
/**
|
/**
|
||||||
@@ -646,7 +481,7 @@ export class StartSdk<Manifest extends T.SDKManifest> {
|
|||||||
effects: E,
|
effects: E,
|
||||||
hostnames: string[],
|
hostnames: string[],
|
||||||
algorithm?: T.Algorithm,
|
algorithm?: T.Algorithm,
|
||||||
) => new GetSslCertificate(effects, hostnames, algorithm),
|
) => new GetSslCertificate(effects, { hostnames, algorithm }),
|
||||||
/** Retrieve the manifest of any installed service package by its ID */
|
/** Retrieve the manifest of any installed service package by its ID */
|
||||||
getServiceManifest,
|
getServiceManifest,
|
||||||
healthCheck: {
|
healthCheck: {
|
||||||
|
|||||||
@@ -1,156 +0,0 @@
|
|||||||
import { Effects } from '../../../base/lib/Effects'
|
|
||||||
import { Manifest, PackageId } from '../../../base/lib/osBindings'
|
|
||||||
import { AbortedError } from '../../../base/lib/util/AbortedError'
|
|
||||||
import { DropGenerator, DropPromise } from '../../../base/lib/util/Drop'
|
|
||||||
import { deepEqual } from '../../../base/lib/util/deepEqual'
|
|
||||||
|
|
||||||
export class GetServiceManifest<Mapped = Manifest> {
|
|
||||||
constructor(
|
|
||||||
readonly effects: Effects,
|
|
||||||
readonly packageId: PackageId,
|
|
||||||
readonly map: (manifest: Manifest | null) => Mapped,
|
|
||||||
readonly eq: (a: Mapped, b: Mapped) => boolean,
|
|
||||||
) {}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Returns the manifest of a service. Reruns the context from which it has been called if the underlying value changes
|
|
||||||
*/
|
|
||||||
async const() {
|
|
||||||
let abort = new AbortController()
|
|
||||||
const watch = this.watch(abort.signal)
|
|
||||||
const res = await watch.next()
|
|
||||||
if (this.effects.constRetry) {
|
|
||||||
watch
|
|
||||||
.next()
|
|
||||||
.then(() => {
|
|
||||||
abort.abort()
|
|
||||||
this.effects.constRetry && this.effects.constRetry()
|
|
||||||
})
|
|
||||||
.catch()
|
|
||||||
}
|
|
||||||
return res.value
|
|
||||||
}
|
|
||||||
/**
|
|
||||||
* Returns the manifest of a service. Does nothing if it changes
|
|
||||||
*/
|
|
||||||
async once() {
|
|
||||||
const manifest = await this.effects.getServiceManifest({
|
|
||||||
packageId: this.packageId,
|
|
||||||
})
|
|
||||||
return this.map(manifest)
|
|
||||||
}
|
|
||||||
|
|
||||||
private async *watchGen(abort?: AbortSignal) {
|
|
||||||
let prev = null as { value: Mapped } | null
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
this.effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
|
||||||
while (this.effects.isInContext && !abort?.aborted) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
const next = this.map(
|
|
||||||
await this.effects.getServiceManifest({
|
|
||||||
packageId: this.packageId,
|
|
||||||
callback: () => callback(),
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
if (!prev || !this.eq(prev.value, next)) {
|
|
||||||
prev = { value: next }
|
|
||||||
yield next
|
|
||||||
}
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
return new Promise<never>((_, rej) => rej(new AbortedError()))
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the manifest of a service. Returns an async iterator that yields whenever the value changes
|
|
||||||
*/
|
|
||||||
watch(abort?: AbortSignal): AsyncGenerator<Mapped, never, unknown> {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
return DropGenerator.of(this.watchGen(ctrl.signal), () => ctrl.abort())
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the manifest of a service. Takes a custom callback function to run whenever it changes
|
|
||||||
*/
|
|
||||||
onChange(
|
|
||||||
callback: (
|
|
||||||
value: Mapped | null,
|
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
|
||||||
) {
|
|
||||||
;(async () => {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of this.watch(ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) {
|
|
||||||
ctrl.abort()
|
|
||||||
break
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetServiceManifest.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetServiceManifest.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the manifest of a service. Returns when the predicate is true
|
|
||||||
*/
|
|
||||||
waitFor(pred: (value: Mapped) => boolean): Promise<Mapped> {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
return DropPromise.of(
|
|
||||||
Promise.resolve().then(async () => {
|
|
||||||
for await (const next of this.watchGen(ctrl.signal)) {
|
|
||||||
if (pred(next)) {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
throw new Error('context left before predicate passed')
|
|
||||||
}),
|
|
||||||
() => ctrl.abort(),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
export function getServiceManifest(
|
|
||||||
effects: Effects,
|
|
||||||
packageId: PackageId,
|
|
||||||
): GetServiceManifest<Manifest>
|
|
||||||
export function getServiceManifest<Mapped>(
|
|
||||||
effects: Effects,
|
|
||||||
packageId: PackageId,
|
|
||||||
map: (manifest: Manifest | null) => Mapped,
|
|
||||||
eq?: (a: Mapped, b: Mapped) => boolean,
|
|
||||||
): GetServiceManifest<Mapped>
|
|
||||||
export function getServiceManifest<Mapped>(
|
|
||||||
effects: Effects,
|
|
||||||
packageId: PackageId,
|
|
||||||
map?: (manifest: Manifest | null) => Mapped,
|
|
||||||
eq?: (a: Mapped, b: Mapped) => boolean,
|
|
||||||
): GetServiceManifest<Mapped> {
|
|
||||||
return new GetServiceManifest(
|
|
||||||
effects,
|
|
||||||
packageId,
|
|
||||||
map ?? ((a) => a as Mapped),
|
|
||||||
eq ?? ((a, b) => deepEqual(a, b)),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
@@ -1,122 +0,0 @@
|
|||||||
import { T } from '..'
|
|
||||||
import { Effects } from '../../../base/lib/Effects'
|
|
||||||
import { AbortedError } from '../../../base/lib/util/AbortedError'
|
|
||||||
import { DropGenerator, DropPromise } from '../../../base/lib/util/Drop'
|
|
||||||
|
|
||||||
export class GetSslCertificate {
|
|
||||||
constructor(
|
|
||||||
readonly effects: Effects,
|
|
||||||
readonly hostnames: string[],
|
|
||||||
readonly algorithm?: T.Algorithm,
|
|
||||||
) {}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Returns the an SSL Certificate for the given hostnames if permitted. Restarts the service if it changes
|
|
||||||
*/
|
|
||||||
const() {
|
|
||||||
return this.effects.getSslCertificate({
|
|
||||||
hostnames: this.hostnames,
|
|
||||||
algorithm: this.algorithm,
|
|
||||||
callback:
|
|
||||||
this.effects.constRetry &&
|
|
||||||
(() => this.effects.constRetry && this.effects.constRetry()),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
/**
|
|
||||||
* Returns the an SSL Certificate for the given hostnames if permitted. Does nothing if it changes
|
|
||||||
*/
|
|
||||||
once() {
|
|
||||||
return this.effects.getSslCertificate({
|
|
||||||
hostnames: this.hostnames,
|
|
||||||
algorithm: this.algorithm,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
private async *watchGen(abort?: AbortSignal) {
|
|
||||||
const resolveCell = { resolve: () => {} }
|
|
||||||
this.effects.onLeaveContext(() => {
|
|
||||||
resolveCell.resolve()
|
|
||||||
})
|
|
||||||
abort?.addEventListener('abort', () => resolveCell.resolve())
|
|
||||||
while (this.effects.isInContext && !abort?.aborted) {
|
|
||||||
let callback: () => void = () => {}
|
|
||||||
const waitForNext = new Promise<void>((resolve) => {
|
|
||||||
callback = resolve
|
|
||||||
resolveCell.resolve = resolve
|
|
||||||
})
|
|
||||||
yield await this.effects.getSslCertificate({
|
|
||||||
hostnames: this.hostnames,
|
|
||||||
algorithm: this.algorithm,
|
|
||||||
callback: () => callback(),
|
|
||||||
})
|
|
||||||
await waitForNext
|
|
||||||
}
|
|
||||||
return new Promise<never>((_, rej) => rej(new AbortedError()))
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the SSL Certificate for the given hostnames if permitted. Returns an async iterator that yields whenever the value changes
|
|
||||||
*/
|
|
||||||
watch(
|
|
||||||
abort?: AbortSignal,
|
|
||||||
): AsyncGenerator<[string, string, string], never, unknown> {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
return DropGenerator.of(this.watchGen(ctrl.signal), () => ctrl.abort())
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the SSL Certificate for the given hostnames if permitted. Takes a custom callback function to run whenever it changes
|
|
||||||
*/
|
|
||||||
onChange(
|
|
||||||
callback: (
|
|
||||||
value: [string, string, string] | null,
|
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
|
||||||
) {
|
|
||||||
;(async () => {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of this.watch(ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) {
|
|
||||||
ctrl.abort()
|
|
||||||
break
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetSslCertificate.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ GetSslCertificate.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Watches the SSL Certificate for the given hostnames if permitted. Returns when the predicate is true
|
|
||||||
*/
|
|
||||||
waitFor(
|
|
||||||
pred: (value: [string, string, string] | null) => boolean,
|
|
||||||
): Promise<[string, string, string] | null> {
|
|
||||||
const ctrl = new AbortController()
|
|
||||||
return DropPromise.of(
|
|
||||||
Promise.resolve().then(async () => {
|
|
||||||
for await (const next of this.watchGen(ctrl.signal)) {
|
|
||||||
if (pred(next)) {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return null
|
|
||||||
}),
|
|
||||||
() => ctrl.abort(),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -4,8 +4,8 @@ import * as TOML from '@iarna/toml'
|
|||||||
import * as INI from 'ini'
|
import * as INI from 'ini'
|
||||||
import * as T from '../../../base/lib/types'
|
import * as T from '../../../base/lib/types'
|
||||||
import * as fs from 'node:fs/promises'
|
import * as fs from 'node:fs/promises'
|
||||||
import { AbortedError, asError, deepEqual } from '../../../base/lib/util'
|
import { asError, deepEqual } from '../../../base/lib/util'
|
||||||
import { DropGenerator, DropPromise } from '../../../base/lib/util/Drop'
|
import { Watchable } from '../../../base/lib/util/Watchable'
|
||||||
import { PathBase } from './Volume'
|
import { PathBase } from './Volume'
|
||||||
|
|
||||||
const previousPath = /(.+?)\/([^/]*)$/
|
const previousPath = /(.+?)\/([^/]*)$/
|
||||||
@@ -228,132 +228,72 @@ export class FileHelper<A> {
|
|||||||
return map(this.validate(data))
|
return map(this.validate(data))
|
||||||
}
|
}
|
||||||
|
|
||||||
private async readConst<B>(
|
private createFileWatchable<B>(
|
||||||
effects: T.Effects,
|
effects: T.Effects,
|
||||||
map: (value: A) => B,
|
map: (value: A) => B,
|
||||||
eq: (left: B | null | undefined, right: B | null) => boolean,
|
eq: (left: B | null, right: B | null) => boolean,
|
||||||
): Promise<B | null> {
|
|
||||||
const watch = this.readWatch(effects, map, eq)
|
|
||||||
const res = await watch.next()
|
|
||||||
if (effects.constRetry) {
|
|
||||||
const record: (typeof this.consts)[number] = [
|
|
||||||
effects.constRetry,
|
|
||||||
res.value,
|
|
||||||
map,
|
|
||||||
eq,
|
|
||||||
]
|
|
||||||
this.consts.push(record)
|
|
||||||
watch
|
|
||||||
.next()
|
|
||||||
.then(() => {
|
|
||||||
this.consts = this.consts.filter((r) => r !== record)
|
|
||||||
effects.constRetry && effects.constRetry()
|
|
||||||
})
|
|
||||||
.catch()
|
|
||||||
}
|
|
||||||
return res.value
|
|
||||||
}
|
|
||||||
|
|
||||||
private async *readWatch<B>(
|
|
||||||
effects: T.Effects,
|
|
||||||
map: (value: A) => B,
|
|
||||||
eq: (left: B | null | undefined, right: B | null) => boolean,
|
|
||||||
abort?: AbortSignal,
|
|
||||||
) {
|
) {
|
||||||
let prev: { value: B | null } | null = null
|
const doRead = async (): Promise<A | null> => {
|
||||||
while (effects.isInContext && !abort?.aborted) {
|
const data = await this.readFile()
|
||||||
if (await exists(this.path)) {
|
if (!data) return null
|
||||||
const ctrl = new AbortController()
|
return this.validate(data)
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
|
||||||
const watch = fs.watch(this.path, {
|
|
||||||
persistent: false,
|
|
||||||
signal: ctrl.signal,
|
|
||||||
})
|
|
||||||
const newRes = await this.readOnce(map)
|
|
||||||
const listen = Promise.resolve()
|
|
||||||
.then(async () => {
|
|
||||||
for await (const _ of watch) {
|
|
||||||
ctrl.abort()
|
|
||||||
return null
|
|
||||||
}
|
|
||||||
})
|
|
||||||
.catch((e) => console.error(asError(e)))
|
|
||||||
if (!prev || !eq(prev.value, newRes)) {
|
|
||||||
console.error('yielding', JSON.stringify({ prev: prev, newRes }))
|
|
||||||
yield newRes
|
|
||||||
}
|
|
||||||
prev = { value: newRes }
|
|
||||||
await listen
|
|
||||||
} else {
|
|
||||||
yield null
|
|
||||||
await onCreated(this.path).catch((e) => console.error(asError(e)))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
return new Promise<never>((_, rej) => rej(new AbortedError()))
|
const filePath = this.path
|
||||||
}
|
const fileHelper = this
|
||||||
|
|
||||||
private readOnChange<B>(
|
const wrappedMap = (raw: A | null): B | null => {
|
||||||
effects: T.Effects,
|
if (raw === null) return null
|
||||||
callback: (
|
return map(raw)
|
||||||
value: B | null,
|
}
|
||||||
error?: Error,
|
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
return new (class extends Watchable<A | null, B | null> {
|
||||||
map: (value: A) => B,
|
protected readonly label = 'FileHelper'
|
||||||
eq: (left: B | null | undefined, right: B | null) => boolean,
|
|
||||||
) {
|
protected async fetch() {
|
||||||
;(async () => {
|
return doRead()
|
||||||
const ctrl = new AbortController()
|
|
||||||
for await (const value of this.readWatch(effects, map, eq, ctrl.signal)) {
|
|
||||||
try {
|
|
||||||
const res = await callback(value)
|
|
||||||
if (res.cancel) ctrl.abort()
|
|
||||||
} catch (e) {
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ FileHelper.read.onChange',
|
|
||||||
e,
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
})()
|
|
||||||
.catch((e) => callback(null, e))
|
|
||||||
.catch((e) =>
|
|
||||||
console.error(
|
|
||||||
'callback function threw an error @ FileHelper.read.onChange',
|
|
||||||
e,
|
|
||||||
),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
|
|
||||||
private readWaitFor<B>(
|
protected async *produce(
|
||||||
effects: T.Effects,
|
abort: AbortSignal,
|
||||||
pred: (value: B | null, error?: Error) => boolean,
|
): AsyncGenerator<A | null, void> {
|
||||||
map: (value: A) => B,
|
while (this.effects.isInContext && !abort.aborted) {
|
||||||
): Promise<B | null> {
|
if (await exists(filePath)) {
|
||||||
const ctrl = new AbortController()
|
const ctrl = new AbortController()
|
||||||
return DropPromise.of(
|
abort.addEventListener('abort', () => ctrl.abort())
|
||||||
Promise.resolve().then(async () => {
|
const watch = fs.watch(filePath, {
|
||||||
const watch = this.readWatch(effects, map, (_) => false, ctrl.signal)
|
persistent: false,
|
||||||
while (true) {
|
signal: ctrl.signal,
|
||||||
try {
|
})
|
||||||
const res = await watch.next()
|
yield await doRead()
|
||||||
if (pred(res.value)) {
|
await Promise.resolve()
|
||||||
ctrl.abort()
|
.then(async () => {
|
||||||
return res.value
|
for await (const _ of watch) {
|
||||||
}
|
ctrl.abort()
|
||||||
if (res.done) {
|
return null
|
||||||
break
|
}
|
||||||
}
|
})
|
||||||
} catch (e) {
|
.catch((e) => console.error(asError(e)))
|
||||||
if (pred(null, e as Error)) {
|
} else {
|
||||||
break
|
yield null
|
||||||
}
|
await onCreated(filePath).catch((e) => console.error(asError(e)))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
ctrl.abort()
|
}
|
||||||
return null
|
|
||||||
}),
|
protected onConstRegistered(value: B | null): (() => void) | void {
|
||||||
() => ctrl.abort(),
|
if (!this.effects.constRetry) return
|
||||||
)
|
const record: (typeof fileHelper.consts)[number] = [
|
||||||
|
this.effects.constRetry,
|
||||||
|
value,
|
||||||
|
wrappedMap,
|
||||||
|
eq,
|
||||||
|
]
|
||||||
|
fileHelper.consts.push(record)
|
||||||
|
return () => {
|
||||||
|
fileHelper.consts = fileHelper.consts.filter((r) => r !== record)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})(effects, { map: wrappedMap, eq })
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -372,7 +312,7 @@ export class FileHelper<A> {
|
|||||||
read(): ReadType<A>
|
read(): ReadType<A>
|
||||||
read<B>(
|
read<B>(
|
||||||
map: (value: A) => B,
|
map: (value: A) => B,
|
||||||
eq?: (left: B | null | undefined, right: B | null) => boolean,
|
eq?: (left: B | null, right: B | null) => boolean,
|
||||||
): ReadType<B>
|
): ReadType<B>
|
||||||
read(
|
read(
|
||||||
map?: (value: A) => any,
|
map?: (value: A) => any,
|
||||||
@@ -382,24 +322,19 @@ export class FileHelper<A> {
|
|||||||
eq = eq ?? deepEqual
|
eq = eq ?? deepEqual
|
||||||
return {
|
return {
|
||||||
once: () => this.readOnce(map),
|
once: () => this.readOnce(map),
|
||||||
const: (effects: T.Effects) => this.readConst(effects, map, eq),
|
const: (effects: T.Effects) =>
|
||||||
watch: (effects: T.Effects, abort?: AbortSignal) => {
|
this.createFileWatchable(effects, map, eq).const(),
|
||||||
const ctrl = new AbortController()
|
watch: (effects: T.Effects, abort?: AbortSignal) =>
|
||||||
abort?.addEventListener('abort', () => ctrl.abort())
|
this.createFileWatchable(effects, map, eq).watch(abort),
|
||||||
return DropGenerator.of(
|
|
||||||
this.readWatch(effects, map, eq, ctrl.signal),
|
|
||||||
() => ctrl.abort(),
|
|
||||||
)
|
|
||||||
},
|
|
||||||
onChange: (
|
onChange: (
|
||||||
effects: T.Effects,
|
effects: T.Effects,
|
||||||
callback: (
|
callback: (
|
||||||
value: A | null,
|
value: A | null,
|
||||||
error?: Error,
|
error?: Error,
|
||||||
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
) => { cancel: boolean } | Promise<{ cancel: boolean }>,
|
||||||
) => this.readOnChange(effects, callback, map, eq),
|
) => this.createFileWatchable(effects, map, eq).onChange(callback),
|
||||||
waitFor: (effects: T.Effects, pred: (value: A | null) => boolean) =>
|
waitFor: (effects: T.Effects, pred: (value: A | null) => boolean) =>
|
||||||
this.readWaitFor(effects, pred, map),
|
this.createFileWatchable(effects, map, eq).waitFor(pred),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,4 @@
|
|||||||
export * from '../../../base/lib/util'
|
export * from '../../../base/lib/util'
|
||||||
export { GetSslCertificate } from './GetSslCertificate'
|
|
||||||
export { GetServiceManifest, getServiceManifest } from './GetServiceManifest'
|
|
||||||
|
|
||||||
export { Drop } from '../../../base/lib/util/Drop'
|
export { Drop } from '../../../base/lib/util/Drop'
|
||||||
export { Volume, Volumes } from './Volume'
|
export { Volume, Volumes } from './Volume'
|
||||||
|
|||||||
Reference in New Issue
Block a user