mirror of
https://github.com/hcengineering/platform.git
synced 2026-08-26 22:32:23 +02:00
+1








98652c6476
* Add bump-changes * Add utility tests * Add utility tests * Bump to new version of esbuild and typescript * v0.7.3 * use platform rig 0.7.10 * upgrade: memory engine optimized; change name to (was recommended by Copilot and Onnikov, TODO: CHANGE CLIENT TOO!!!) Signed-off-by: Leonid Kaganov <lleo@lleo.me> * Fix rate limits bug * Bump versions * Fix lock file * Fix bug in queue cleanup * Add more tests for queue * Add api-test tests * Initial commit * Improve hierarchy + tests Add tests for hierarchy and few performance/memory optimizations. * Add more hierarchy tests * Move from Huly platform repository * Add docker tests setup * Fix test to be executed only once * Add connection tests * Fix package include source files * More tests * Create README.md * Fix pnpm lock * Fix packages publish * Remove broken tests * feat: adjust hulylake client for storage adapter Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Fix export * Fix publish * Fix message update (#114) Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Bump version Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Bump versions Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Fix lang store (#115) Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * update hulylake client Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump version Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Add hulylake storage adapter Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump version Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fix validation issues Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fix: do not fail on deseralization error and add logs Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * bump version -> 0.1.14 Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fix collaboration test Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Update prettier and new update-deps script Prettier + svelte support * fix unstable ydoc tests Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Add tx ordering middleware * Fix ordering tests * Fix Kafka close of admin * Add tests for measurement and understand overhead * Fix not updated lock file * Fix update-deps * Fix update-deps * Use latest platform-rig * Fix deps * Add rush check to CI * Use latest versions * Bump versions * Fix lock file * validate json patch Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * bump version -> 0.1.15 Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fix merge unit tests Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Script to sync eslint deps * Fix deps * Fix tests * Fix platform-rig detection * Update to latest platform-rig * Update to latest platform rig and core * Bump typescript * Bump typescript * Rollback eslint plugins * Fix lock file * Bump platform-rig * Update to latest platform-rig * update to latest platform-rig * Allow to compile svelte files * Add ui-test component for checking compile * Fix log levels rename compile ui -> compile ui-esbuild * Fix build * Bump esbuild svelte version * Chore: use fixed versions in update-deps Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Chore: commit changes Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Update deps * Add tests for session manager * Fix txOrdering implementation * Bump ordering * Prevent metrics zero values in measure + Fix format svelte files * Revert update-deps script logic * v0.7.19 * update to latest platform-rig * Update deps * Fix pnpm * Session counters * Fix pnpm lock * Add storage client Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump versions Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fix versions Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump core * Fix pnpm * Get rid of communication dependency * Add copilot memory file * Use proper name for instructions file * Fix instructions * Use domain instead of test name in gauges * Update instructions file * fix front service upload Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * remove incorrect test Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Move packages to huly.core * Move packages to core, since they are not utils * Add global user profile Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Fix lock file Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Add support for memory limit check * Bump version * Fix pnpm * report more accurate upload progress Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fic validation issues Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Fix deps Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Move LowLevelStorage to server * Fix linting * Revert "Fix linting" This reverts commit54631d353e. * Revert "Move LowLevelStorage to server" This reverts commitaafb8f6f12. * feature: add regorus engine with permit file Signed-off-by: Leonid Kaganov <lleo@lleo.me> * Fix one second counters for memory usage * Fix kafka test * use fresh core * Version bump * fix: key parameter added Signed-off-by: Leonid Kaganov <lleo@lleo.me> * Fix readme and few author mistakes * Export domain schemas * Bump version * Tests (#117) Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Bump version Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * feat: compact compact worker (#4) Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * bump version -> 0.1.16 Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Add TypeIdentifier Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Add change logs Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * rename send -> try_send Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Fix pnpm lock Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Add identifier middleware, bump core Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Add subsciption methods to account client Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Fix lock file Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Fix reaction notification (#118) Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Bump version Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Improve find methods schemas to convert to valid types Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> * Add change description Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> * Do not transcode while recording Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Open telemetry support Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * use proper content type in multipart upload Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump versions Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * fix build (#26) Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Fix peers (#120) Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Bump version Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> * Add ActivityCollaborativeChange Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Update pnpm Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Allow to suspend errors on with * Fix pnpm cache * update versions * v0.7.17 for all * v0.7.11 * v0.7.14 * remove arc from worker Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * bump version -> 0.1.17 Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Rank for attributes Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Update pnpm Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Fix one second counters * Fix withContext and allow pass options * Fix formatting * Use updated deps * Bump versions * Update deps * Update deps to platform.core * add support for textColor mark Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * add support for textStyle mark Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump versions Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump versions again Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Fix Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Rework on second timers * fix merge of large blobs feched from s3 Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * bump version -> 0.1.18 Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * New subscription methods in account-client * Update lock file Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Send error on find for wrong domain * Suspend connect custom errors events in traces * Bump client * Bump core * update deps * Fix lock file * Sorting for TypeIdentifier Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Bump version Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * add workspace usage info Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Bump versions Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> * Improve pg security perfomance Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Fix identifier middleware Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Update TxAccessLevel interface Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Allow guest to update its identities Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Add password login locked platform status Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> * Fix Uptrace normalizeMarkdown errors Signed-off-by: Artem Savchenko <armisav@gmail.com> * Add change log Signed-off-by: Artem Savchenko <armisav@gmail.com> * Add txMatch to permission Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * update pnpm lock Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Bump Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Fix permission middleware Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Fix enum sorting Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Enable formatting check Signed-off-by: Andrey Sobolev <haiodo@gmail.com> * Enable formatting check * Add change Signed-off-by: Andrey Sobolev <haiodo@gmail.com> * Fix Uptrace NaN error Signed-off-by: Artem Savchenko <armisav@gmail.com> * feature: removed actors, improved performance Signed-off-by: Leonid Kaganov <lleo@lleo.me> * feature: ping from server to clients added Signed-off-by: Leonid Kaganov <lleo@lleo.me> * feature: ping from server to clients added Signed-off-by: Leonid Kaganov <lleo@lleo.me> * Compress kafka messages and fix exception in findAll Signed-off-by: Artem Savchenko <armisav@gmail.com> * Bump versions Signed-off-by: Artem Savchenko <armisav@gmail.com> * Bump versions Signed-off-by: Artem Savchenko <armisav@gmail.com> * Rush change Signed-off-by: Artem Savchenko <armisav@gmail.com> * Fix compression param Signed-off-by: Artem Savchenko <armisav@gmail.com> * Trigger change Signed-off-by: Artem Savchenko <armisav@gmail.com> * Clean up Signed-off-by: Artem Savchenko <armisav@gmail.com> * Trigger change Signed-off-by: Artem Savchenko <armisav@gmail.com> * Bump markdown version Signed-off-by: Artem Savchenko <armisav@gmail.com> * Enable sub projects * Fix wrong double symbol scripts * Include foundation packages * Add support for custom exclude filters Add support for custom exclude filters - by Andrey Sobolev - haiodo@gmail.com Signed-off-by: Andrey Sobolev <haiodo@gmail.com> * Bump Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> * Fix Uptrace filter is not a function error Signed-off-by: Artem Savchenko <armisav@gmail.com> * Sync versions Signed-off-by: Andrey Sobolev <haiodo@gmail.com> --------- Signed-off-by: Leonid Kaganov <lleo@lleo.me> Signed-off-by: Alexander Onnikov <Alexander.Onnikov@xored.com> Signed-off-by: Kristina Fefelova <kristin.fefelova@gmail.com> Signed-off-by: Alexey Zinoviev <alexey.zinoviev@xored.com> Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com> Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> Signed-off-by: Artem Savchenko <armisav@gmail.com> Signed-off-by: Andrey Sobolev <haiodo@gmail.com> Co-authored-by: Leonid Kaganov <lleo@lleo.me> Co-authored-by: Alexander Onnikov <Alexander.Onnikov@xored.com> Co-authored-by: Alexander Onnikov <Alexander.Onnikov@gmail.com> Co-authored-by: Kristina <kristin.fefelova@gmail.com> Co-authored-by: Alexey Zinoviev <alexey.zinoviev@xored.com> Co-authored-by: Denis Bykhov <bykhov.denis@gmail.com> Co-authored-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> Co-authored-by: Alexander Onnikov <aonnikov@hardcoreeng.com> Co-authored-by: Artem Savchenko <armisav@gmail.com>
463 lines
14 KiB
TypeScript
463 lines
14 KiB
TypeScript
import { BackRPCClient, type BackRPCResponseSend } from '@hcengineering/network-backrpc'
|
|
import {
|
|
agentDirectRef,
|
|
EndpointKind,
|
|
parseEndpointRef,
|
|
type AgentEndpointRef,
|
|
type AgentRecordInfo,
|
|
type AgentUuid,
|
|
type ClientUuid,
|
|
type ContainerConnection,
|
|
type ContainerEndpointRef,
|
|
type NetworkEvent,
|
|
type ContainerKind,
|
|
type ContainerRecord,
|
|
type ContainerReference,
|
|
type NetworkUpdateListener,
|
|
type ContainerUuid,
|
|
type GetOptions,
|
|
type NetworkAgent,
|
|
type NetworkClient,
|
|
type TickManager,
|
|
NetworkEventKind,
|
|
createProxy
|
|
} from '@hcengineering/network-core'
|
|
import { v4 as uuidv4 } from 'uuid'
|
|
import { ContainerConnectionImpl, NetworkDirectConnectionImpl, RoutedNetworkAgentConnectionImpl } from './agent'
|
|
import { opNames } from './types'
|
|
|
|
interface ClientAgentRecord {
|
|
agent: NetworkAgent
|
|
register: Promise<void>
|
|
resolve: () => void
|
|
}
|
|
|
|
class ContainerReferenceImpl implements ContainerReference {
|
|
constructor (
|
|
readonly uuid: ContainerUuid,
|
|
private readonly client: NetworkClientImpl
|
|
) {}
|
|
|
|
get endpoint (): ContainerEndpointRef {
|
|
const ref = this.client.references.get(this.uuid)
|
|
if (ref === undefined) {
|
|
throw new Error('Reference not found')
|
|
}
|
|
return ref.endpoint
|
|
}
|
|
|
|
async close (): Promise<void> {
|
|
await this.client.release(this.uuid)
|
|
this.client.references.delete(this.uuid)
|
|
}
|
|
|
|
async request (operation: string, data?: any): Promise<any> {
|
|
return await this.client.request(this.uuid, operation, data)
|
|
}
|
|
|
|
cast<T extends object>(interfaceName?: string): T {
|
|
return createProxy<T>(this, interfaceName)
|
|
}
|
|
|
|
async connect (): Promise<ContainerConnection> {
|
|
let conn = this.client.containerConnections.get(this.uuid)
|
|
if (conn !== undefined) {
|
|
return conn
|
|
}
|
|
conn = this.client.establishConnection(this.uuid, this.endpoint)
|
|
await conn.connect()
|
|
return conn
|
|
}
|
|
}
|
|
|
|
interface ContainerRef {
|
|
ref: ContainerReference
|
|
kind: ContainerKind
|
|
request: GetOptions
|
|
endpoint: ContainerEndpointRef
|
|
}
|
|
|
|
/**
|
|
* Huly Network client
|
|
*
|
|
* Some methods are omit clientId parameter.
|
|
*/
|
|
export class NetworkClientImpl implements NetworkClient {
|
|
clientId: ClientUuid = uuidv4() as ClientUuid
|
|
|
|
private readonly client: BackRPCClient<ClientUuid>
|
|
|
|
readonly _agents = new Map<AgentUuid, ClientAgentRecord>()
|
|
|
|
// A set of clients for individual containers or agent TORs
|
|
containerConnections = new Map<ContainerUuid, ContainerConnectionImpl>()
|
|
agentConnections = new Map<AgentEndpointRef, RoutedNetworkAgentConnectionImpl<ClientUuid>>()
|
|
|
|
cid: number = 0
|
|
containerListeners = new Map<number, NetworkUpdateListener>()
|
|
|
|
references = new Map<ContainerUuid, ContainerRef>()
|
|
|
|
registered: boolean = false
|
|
|
|
constructor (
|
|
readonly host: string,
|
|
port: number,
|
|
protected readonly tickMgr: TickManager,
|
|
aliveTimeout?: number
|
|
) {
|
|
const options = undefined
|
|
this.client = new BackRPCClient<ClientUuid>(this.clientId, this, host, port, tickMgr, options, aliveTimeout)
|
|
}
|
|
|
|
async waitConnection (timeout?: number): Promise<void> {
|
|
if (timeout !== undefined) {
|
|
await new Promise<void>((resolve, reject) => {
|
|
const co = setTimeout(() => {
|
|
// Timeout reached, we reject the promise by throwing an error
|
|
reject(new Error('Connection timeout'))
|
|
}, timeout)
|
|
|
|
this.client
|
|
.waitConnection()
|
|
.then(() => {
|
|
resolve()
|
|
clearTimeout(co)
|
|
})
|
|
.catch((err) => {
|
|
reject(err)
|
|
})
|
|
})
|
|
return
|
|
}
|
|
await this.client.waitConnection()
|
|
}
|
|
|
|
async close (): Promise<void> {
|
|
for (const refs of this.references.values()) {
|
|
await refs.ref.close()
|
|
}
|
|
for (const directConn of this.containerConnections.values()) {
|
|
await directConn.close()
|
|
}
|
|
for (const agentConn of this.agentConnections.values()) {
|
|
await agentConn.close()
|
|
}
|
|
|
|
for (const agent of this._agents.values()) {
|
|
await this.client.request<ContainerEndpointRef[]>(opNames.unregister, {
|
|
uuid: agent.agent.uuid
|
|
})
|
|
}
|
|
this.client.close()
|
|
}
|
|
|
|
async requestHandler (method: string, params: any, send: BackRPCResponseSend): Promise<void> {
|
|
const [agentId, agentParams] = params
|
|
// Pass agent methods to a proper agent
|
|
const { agent } = this._agents.get(agentId) ?? { agent: undefined }
|
|
if (agent === undefined) {
|
|
await send({ error: `Agent ${agentId} not found` })
|
|
return
|
|
}
|
|
switch (method) {
|
|
case opNames.getContainer:
|
|
await send(await agent.get(agentParams[0], agentParams[1]))
|
|
break
|
|
case opNames.listContainers:
|
|
await send(await agent.list(agentParams[0]))
|
|
break
|
|
case opNames.sendContainer:
|
|
await send(await agent.request(agentParams[0], agentParams[1], agentParams[2]))
|
|
break
|
|
case opNames.terminate:
|
|
await agent.terminate(agentParams[0] as ContainerUuid)
|
|
await send('')
|
|
break
|
|
default:
|
|
throw new Error('Unknown method')
|
|
}
|
|
}
|
|
|
|
async onEvent (event: NetworkEvent): Promise<void> {
|
|
// Handle container events
|
|
// In case of container stopped, agent stopped or endpoint changed, we need to update direct connections to be re-established.
|
|
await this.handleConnectionUpdates(event)
|
|
|
|
// Handle container removal for stateless containers - attempt to re-register
|
|
for (const containerEvent of event.containers) {
|
|
if (containerEvent.event === NetworkEventKind.removed) {
|
|
// Check if any of our agents have this container as stateless and need to re-register
|
|
for (const agentRecord of this._agents.values()) {
|
|
const agent = agentRecord.agent as any
|
|
const statelessContainers = agent.statelessContainers as Map<ContainerUuid, any> | undefined
|
|
if (statelessContainers !== undefined && statelessContainers.has(containerEvent.container.uuid)) {
|
|
console.log(
|
|
`HA: Container ${containerEvent.container.uuid} removed, attempting to re-register from agent ${agent.uuid}`
|
|
)
|
|
// Re-register this agent to attempt to claim the container
|
|
setTimeout(() => {
|
|
this.doRegister(agent).catch((err) => {
|
|
console.error(`Failed to re-register agent ${agent.uuid}:`, err)
|
|
})
|
|
}, 100) // Small delay to avoid thundering herd
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
for (const listener of this.containerListeners.values()) {
|
|
try {
|
|
await listener(event)
|
|
} catch (error) {
|
|
console.error('Error in container listener:', error)
|
|
}
|
|
}
|
|
}
|
|
|
|
async handleRefUpdate (uuid: ContainerUuid, endpoint: ContainerEndpointRef): Promise<void> {
|
|
const ref = this.references.get(uuid)
|
|
if (ref !== undefined) {
|
|
const conn = this.containerConnections.get(ref.ref.uuid)
|
|
if (conn !== undefined && ref.endpoint !== endpoint) {
|
|
conn.setConnection(this.establishConnection(ref.ref.uuid, endpoint))
|
|
} else {
|
|
ref.endpoint = endpoint
|
|
}
|
|
}
|
|
}
|
|
|
|
async handleNewContainer (oldUuid: ContainerUuid, uuid: ContainerUuid, endpoint: ContainerEndpointRef): Promise<void> {
|
|
const ref = this.references.get(oldUuid)
|
|
this.references.delete(oldUuid)
|
|
if (ref !== undefined) {
|
|
const conn = this.containerConnections.get(oldUuid)
|
|
this.containerConnections.delete(oldUuid)
|
|
if (conn !== undefined) {
|
|
this.containerConnections.set(uuid, conn)
|
|
if (ref.endpoint !== endpoint) {
|
|
conn.setConnection(this.establishConnection(uuid, endpoint))
|
|
}
|
|
}
|
|
ref.ref.uuid = uuid
|
|
ref.endpoint = endpoint
|
|
|
|
this.references.set(uuid, ref)
|
|
}
|
|
}
|
|
|
|
async onRegister (): Promise<void> {
|
|
this.registered = true
|
|
// We need to re-register all our managed agents, since we could provide containers we request to our selfs
|
|
for (const agent of this._agents.values()) {
|
|
await this.doRegister(agent.agent)
|
|
}
|
|
|
|
for (const [uuid, ref] of this.references.entries()) {
|
|
const [newUuid, newEndpoint] = await this.retryGetContainerRef(ref.kind, ref.request)
|
|
if (uuid !== newUuid) {
|
|
await this.handleNewContainer(uuid, newUuid, newEndpoint)
|
|
} else {
|
|
await this.handleRefUpdate(uuid, newEndpoint)
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Register a new agent, agent could or could not provide an endpoint for routed connections.
|
|
*/
|
|
async register (agent: NetworkAgent): Promise<void> {
|
|
const rec: ClientAgentRecord = {
|
|
agent,
|
|
register: Promise.resolve(),
|
|
resolve: () => {}
|
|
}
|
|
rec.register = new Promise<void>((resolve) => {
|
|
rec.resolve = resolve
|
|
})
|
|
this._agents.set(agent.uuid, rec)
|
|
|
|
agent.onUpdate = async (event) => {
|
|
await this.client.request(opNames.containerUpdate, event)
|
|
}
|
|
agent.onAgentUpdate = async () => {
|
|
await this.doRegister(agent)
|
|
}
|
|
|
|
if (this.registered) {
|
|
await this.doRegister(agent)
|
|
}
|
|
await rec.register
|
|
}
|
|
|
|
async doRegister (agent: NetworkAgent): Promise<void> {
|
|
const containers: ContainerRecord[] = []
|
|
for (const container of await agent.list()) {
|
|
containers.push({
|
|
agentId: agent.uuid,
|
|
uuid: container.uuid,
|
|
endpoint: container.endpoint,
|
|
kind: container.kind,
|
|
lastVisit: container.lastVisit
|
|
} satisfies ContainerRecord)
|
|
}
|
|
const toClean = await this.client.request<ContainerUuid[]>(opNames.register, {
|
|
uuid: agent.uuid,
|
|
containers,
|
|
kinds: agent.kinds,
|
|
endpoint: agent.endpoint
|
|
})
|
|
for (const uuid of toClean) {
|
|
await agent.terminate(uuid)
|
|
}
|
|
this._agents.get(agent.uuid)?.resolve()
|
|
}
|
|
|
|
async agents (): Promise<AgentRecordInfo[]> {
|
|
// Return actual list of agents
|
|
return await this.client.request<AgentRecordInfo[]>(opNames.getAgents, {})
|
|
}
|
|
|
|
async kinds (): Promise<ContainerKind[]> {
|
|
return await this.client.request<ContainerKind[]>(opNames.getKinds, {})
|
|
}
|
|
|
|
async get (kind: ContainerKind, request: GetOptions): Promise<ContainerReference> {
|
|
// TODO: Wait for all pending requests to finish
|
|
|
|
if (request.uuid !== undefined) {
|
|
const existing = this.references.get(request.uuid)
|
|
if (existing !== undefined) {
|
|
return existing.ref
|
|
}
|
|
}
|
|
const [uuid, endpoint] = await this.retryGetContainerRef(kind, request)
|
|
const ref: ContainerReference = new ContainerReferenceImpl(uuid, this)
|
|
this.references.set(uuid, { kind, ref, request, endpoint })
|
|
return ref
|
|
}
|
|
|
|
establishConnection (uuid: ContainerUuid, endpoint: ContainerEndpointRef): ContainerConnectionImpl {
|
|
// Check if connection is routed
|
|
const parsedRef = parseEndpointRef(endpoint)
|
|
if (parsedRef.uuid === undefined) {
|
|
throw new Error('Invalid endpoint reference')
|
|
}
|
|
if (parsedRef.kind === EndpointKind.noconnect) {
|
|
throw new Error('No connection available')
|
|
}
|
|
if (parsedRef.kind === EndpointKind.routed) {
|
|
const agentRef = agentDirectRef(parsedRef.host, parsedRef.port, parsedRef.agentId)
|
|
let agentConn = this.agentConnections.get(agentRef)
|
|
if (agentConn === undefined) {
|
|
agentConn = new RoutedNetworkAgentConnectionImpl<ClientUuid>(
|
|
this.tickMgr,
|
|
this.clientId,
|
|
parsedRef.host,
|
|
parsedRef.port
|
|
)
|
|
this.agentConnections.set(agentRef, agentConn)
|
|
}
|
|
let conn = this.containerConnections.get(uuid)
|
|
if (conn === undefined) {
|
|
conn = new ContainerConnectionImpl(uuid, agentConn.connect(parsedRef.uuid))
|
|
} else {
|
|
conn.setConnection(agentConn.connect(parsedRef.uuid))
|
|
}
|
|
this.containerConnections.set(uuid, conn)
|
|
return conn
|
|
}
|
|
const directConn = new NetworkDirectConnectionImpl(
|
|
this.tickMgr,
|
|
this.clientId,
|
|
parsedRef.uuid,
|
|
parsedRef.host,
|
|
parsedRef.port
|
|
)
|
|
let conn = this.containerConnections.get(uuid)
|
|
if (conn === undefined) {
|
|
conn = new ContainerConnectionImpl(uuid, directConn)
|
|
} else {
|
|
conn.setConnection(directConn)
|
|
}
|
|
this.containerConnections.set(uuid, conn)
|
|
return conn
|
|
}
|
|
|
|
async handleConnectionUpdates (event: NetworkEvent): Promise<void> {
|
|
// Handle connection updates
|
|
for (const e of event.containers ?? []) {
|
|
if (e.event === NetworkEventKind.removed || e.event === NetworkEventKind.updated) {
|
|
await this.handleRefUpdate(e.container.uuid, e.container.endpoint)
|
|
}
|
|
}
|
|
}
|
|
|
|
private async getContainerRef (
|
|
kind: ContainerKind,
|
|
request: GetOptions
|
|
): Promise<[ContainerUuid, ContainerEndpointRef]> {
|
|
return await this.client.request<[ContainerUuid, ContainerEndpointRef]>(opNames.getContainer, { kind, request })
|
|
}
|
|
|
|
async retryGetContainerRef (kind: ContainerKind, request: GetOptions): Promise<[ContainerUuid, ContainerEndpointRef]> {
|
|
let waitTimeout: number = 1
|
|
let earlyRetry = (): void => {}
|
|
const stop = this.onUpdate(async (event) => {
|
|
// We agent is appear with a required kind, we can retry immediately
|
|
if (event.agents.some((e) => e.event === NetworkEventKind.added && e.kinds.includes(kind))) {
|
|
waitTimeout = 0
|
|
earlyRetry()
|
|
}
|
|
})
|
|
try {
|
|
while (true) {
|
|
try {
|
|
const ref = await this.getContainerRef(kind, request)
|
|
if (waitTimeout > 1) {
|
|
console.log(`Successfully got container ref for ${kind} after ${waitTimeout - 1} retries.`)
|
|
}
|
|
return ref
|
|
} catch (err) {
|
|
console.warn(`Error getting container ref for ${kind}. Will retry...`)
|
|
|
|
await Promise.any([
|
|
this.tickMgr.waitTick(waitTimeout),
|
|
new Promise<void>((resolve) => {
|
|
earlyRetry = resolve
|
|
})
|
|
])
|
|
if (waitTimeout < this.tickMgr.tps * 5) {
|
|
waitTimeout++
|
|
}
|
|
}
|
|
}
|
|
} finally {
|
|
stop()
|
|
}
|
|
}
|
|
|
|
async release (uuid: ContainerUuid): Promise<void> {
|
|
await this.client.request<any>(opNames.releaseContainer, { uuid })
|
|
}
|
|
|
|
async list (kind?: ContainerKind): Promise<ContainerRecord[]> {
|
|
return await this.client.request<ContainerRecord[]>(opNames.listContainers, {
|
|
kind
|
|
})
|
|
}
|
|
|
|
// Send some data to container, using proxy connection.
|
|
async request (target: ContainerUuid, operation: string, data?: any): Promise<any> {
|
|
return await this.client.request<any>(opNames.sendContainer, [target, operation, data])
|
|
}
|
|
|
|
onUpdate (listener: NetworkUpdateListener): () => void {
|
|
const cid = this.cid++
|
|
this.containerListeners.set(cid, listener)
|
|
return () => {
|
|
this.containerListeners.delete(cid)
|
|
}
|
|
}
|
|
}
|