diff --git a/.vscode/launch.json b/.vscode/launch.json index 3c237f20d0..c13806bcbb 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -617,6 +617,7 @@ }, "runtimeArgs": ["--nolazy", "-r", "ts-node/register"], "sourceMaps": true, + "nodeVersionHint": 22, "cwd": "${workspaceRoot}/services/github/pod-github", "protocol": "inspector", "outputCapture": "std" diff --git a/dev/tool/src/db.ts b/dev/tool/src/db.ts index a32c875672..37b9c4d672 100644 --- a/dev/tool/src/db.ts +++ b/dev/tool/src/db.ts @@ -1319,9 +1319,41 @@ async function migrateAccount ( accountDB: AccountDB, dryRun = true ): Promise { - const primaryKey: SocialKey = { - type: SocialIdType.EMAIL, - value: account.email + let primaryKey: SocialKey + let secondaryKey: SocialKey | undefined + + if (account.githubId != null) { + if (account.githubUser == null) { + console.log('No github user found for github id', account.githubId) + return + } + + primaryKey = { + type: SocialIdType.GITHUB, + value: account.githubUser + } + secondaryKey = !account.email.startsWith('github:') + ? { + type: SocialIdType.EMAIL, + value: account.email + } + : undefined + } else if (account.openId != null) { + primaryKey = { + type: SocialIdType.OIDC, + value: account.openId + } + secondaryKey = !account.email.startsWith('openid:') + ? { + type: SocialIdType.EMAIL, + value: account.email + } + : undefined + } else { + primaryKey = { + type: SocialIdType.EMAIL, + value: account.email + } } let personUuid: PersonUuid @@ -1389,6 +1421,21 @@ async function migrateAccount ( } } + if (secondaryKey != null) { + const existingSecondary = await accountDB.socialId.findOne(secondaryKey) + if (existingSecondary == null) { + if (!dryRun) { + await accountDB.socialId.insertOne({ + ...secondaryKey, + personUuid, + ...verified + }) + } else { + console.log('Creating secondary social id', { personUuid, confirmed: account.confirmed, secondaryKey }) + } + } + } + return personUuid as AccountUuid } diff --git a/models/ai-assistant/src/index.ts b/models/ai-assistant/src/index.ts index d4d8b64101..8996bcd792 100644 --- a/models/ai-assistant/src/index.ts +++ b/models/ai-assistant/src/index.ts @@ -33,6 +33,7 @@ export function createModel (builder: Builder): void { icon: aiAssistant.component.IconHulyAssistant, allowMultiple: false, createComponent: aiAssistant.component.Connect, + configureComponent: aiAssistant.component.Connect, onDisconnect: aiAssistant.handler.DisconnectHandler, onDisconnectAll: aiAssistant.handler.DisconnectAllHandler, reconnectComponent: aiAssistant.component.Connect, diff --git a/plugins/chat-resources/src/components/ChatNavigation.svelte b/plugins/chat-resources/src/components/ChatNavigation.svelte index f0bdaa5251..34cd5b34d8 100644 --- a/plugins/chat-resources/src/components/ChatNavigation.svelte +++ b/plugins/chat-resources/src/components/ChatNavigation.svelte @@ -116,9 +116,6 @@ { type: communication.type.Direct, order: 2 } ], fixedTypes: [chat.masterTag.Thread, communication.type.Direct], - specialSorting: { - [communication.type.Direct]: 'alphabetical' - }, allowCreate: true, defaultSorting: 'recent', lookback: '2w', diff --git a/plugins/client-resources/src/connection.ts b/plugins/client-resources/src/connection.ts index 456c881ed6..fd4b2bcd08 100644 --- a/plugins/client-resources/src/connection.ts +++ b/plugins/client-resources/src/connection.ts @@ -72,7 +72,7 @@ class RequestPromise { reject!: (reason?: any) => void reconnect?: () => void - // Required to proeprly handle rate limits + // Required to properly handle rate limits sendData: () => void = () => {} constructor ( @@ -92,6 +92,11 @@ class RequestPromise { const globalRPCHandler: RPCHandler = new RPCHandler() +interface OnConnectHandler { + resolve: () => void + reject: (err: Error) => void +} + class Connection implements ClientConnection { private websocket: ClientSocket | null = null binaryMode = false @@ -99,7 +104,7 @@ class Connection implements ClientConnection { private readonly requests = new Map() private lastId = 0 private interval: number | undefined - private dialTimer: any | undefined + private dialTimer: number | undefined private sockets = 0 private openAction: any @@ -166,12 +171,13 @@ class Connection implements ClientConnection { } private schedulePing (socketId: number): void { - clearInterval(this.interval) this.pingResponse = Date.now() const wsocket = this.websocket - const interval = setInterval(() => { + + clearInterval(this.interval) + this.interval = setInterval(() => { if (wsocket !== this.websocket) { - clearInterval(interval) + clearInterval(this.interval) return } if (!this.upgrading && this.pingResponse !== 0 && Date.now() - this.pingResponse > hangTimeout) { @@ -203,7 +209,6 @@ class Connection implements ClientConnection { clearInterval(this.interval) } }, pingTimeout) - this.interval = interval } async close (): Promise { @@ -211,6 +216,12 @@ class Connection implements ClientConnection { clearTimeout(this.openAction) clearTimeout(this.dialTimer) clearInterval(this.interval) + for (const handler of this.onConnectHandlers) { + handler.reject(new Error('Connection closed')) + } + for (const req of this.requests.values()) { + req.reject(new Error('Connection closed')) + } if (this.websocket !== null) { this.websocket.close(1000) this.websocket = null @@ -222,7 +233,7 @@ class Connection implements ClientConnection { } delay = 0 - onConnectHandlers: (() => void)[] = [] + onConnectHandlers: OnConnectHandler[] = [] private waitOpenConnection (ctx: MeasureContext): Promise | undefined { if (this.isConnected()) { @@ -233,9 +244,10 @@ class Connection implements ClientConnection { 'wait-connection', {}, (ctx) => - new Promise((resolve) => { - this.onConnectHandlers.push(() => { - resolve() + new Promise((resolve, reject) => { + this.onConnectHandlers.push({ + resolve, + reject }) // Websocket is null for first time this.scheduleOpen(ctx, false) @@ -294,7 +306,7 @@ class Connection implements ClientConnection { if (resp.error !== undefined) { if (resp.error?.code === UNAUTHORIZED.code || resp.terminate === true) { if ( - resp.error.code !== platform.status.WorkspaceArchived || + resp.error.code !== platform.status.WorkspaceArchived && resp.error.code !== platform.status.WorkspaceNotFound ) { Analytics.handleError(new PlatformError(resp.error)) @@ -350,7 +362,7 @@ class Connection implements ClientConnection { // We need to clear dial timer, since we recieve hello response. clearTimeout(this.dialTimer) - this.dialTimer = null + this.dialTimer = undefined this.lastHash = (resp as HelloResponse).lastHash const serverVersion = helloResp.serverVersion @@ -374,7 +386,7 @@ class Connection implements ClientConnection { // Notify all waiting connection listeners const handlers = this.onConnectHandlers.splice(0, this.onConnectHandlers.length) for (const h of handlers) { - h() + h.resolve() } for (const [, v] of this.requests.entries()) { @@ -540,12 +552,10 @@ class Connection implements ClientConnection { return } this.websocket = wsocket - const opened = false - - if (this.dialTimer != null) { + if (this.dialTimer === undefined) { this.dialTimer = setTimeout(() => { - this.dialTimer = null - if (!opened && !this.closed) { + this.dialTimer = undefined + if (!this.closed) { void this.opt?.onDialTimeout?.()?.catch((err) => { this.ctx.error('failed to handle dial timeout', { err }) }) @@ -652,16 +662,15 @@ class Connection implements ClientConnection { ctx.withSync('send-hello', {}, () => this.websocket?.send(this.rpcHandler.serialize(helloRequest, false))) } - wsocket.onerror = (event: any) => { + // FIX: remove undefined variable 'opened' + wsocket.onerror = () => { if (this.websocket !== wsocket) { return } if (this.delay < 3) { this.delay += 1 } - if (opened) { - console.error('client websocket error:', socketId, this.url, this.workspace, this.user) - } + console.error('client websocket error:', socketId, this.url, this.workspace, this.user) } } @@ -761,7 +770,7 @@ class Connection implements ClientConnection { getAccount (): Promise { if (this.account !== undefined) { - return clone(this.account) + return Promise.resolve(clone(this.account)) } return this.sendRequest({ method: 'getAccount', params: [] }) } diff --git a/plugins/controlled-documents-resources/src/components/document/EditDocContent.svelte b/plugins/controlled-documents-resources/src/components/document/EditDocContent.svelte index 5c7b4e71fd..ae079c94d2 100644 --- a/plugins/controlled-documents-resources/src/components/document/EditDocContent.svelte +++ b/plugins/controlled-documents-resources/src/components/document/EditDocContent.svelte @@ -25,7 +25,8 @@ TableOfContents, TableOfContentsContent, getNodeElement, - highlightUpdateCommand + highlightUpdateCommand, + selectNode } from '@hcengineering/text-editor-resources' import { EditBox, Label, Scroller } from '@hcengineering/ui' import { getCollaborationUser } from '@hcengineering/view-resources' @@ -41,8 +42,8 @@ $documentCommentHighlightedLocation as documentCommentHighlightedLocation, $documentComments as documentComments, documentCommentsDisplayRequested, - documentCommentsHighlightUpdated, documentCommentsLocationNavigateRequested, + documentCommentsAddCanceled, $isEditable as isEditable } from '../../stores/editors/document' import DocumentPrintTitlePage from '../print/DocumentPrintTitlePage.svelte' @@ -58,19 +59,16 @@ let headings: Heading[] = [] let textEditor: CollaboratorEditor let selectedNodeId: string | null | undefined = undefined - let isFocused = false let editor: Editor let title = $controlledDocument?.title ?? '' $: isTemplate = $controlledDocument != null && hierarchy.hasMixin($controlledDocument, documents.mixin.DocumentTemplate) - function handleRefreshHighlight () { - if (!textEditor) { - return - } + $: commentUuids = $documentComments.map((p) => p.nodeId).filter((id) => id != null) - textEditor.commands()?.command(highlightUpdateCommand()) + function handleRefreshHighlight (): void { + textEditor?.commands()?.command(highlightUpdateCommand()) } const unsubscribeHighlightRefresh = merge([documentCommentHighlightedLocation, documentComments.updates]).subscribe({ @@ -82,21 +80,28 @@ const unsubscribeNavigateToLocation = documentCommentsLocationNavigateRequested.subscribe({ // eslint-disable-next-line @typescript-eslint/no-misused-promises next: async ({ nodeId }) => { - if (!nodeId) { + if (nodeId == null) { handleRefreshHighlight() return } - if (!textEditor) { - return + selectedNodeId = nodeId + + if (editor !== undefined) { + await tick() + + const element = getNodeElement(editor, nodeId) + element?.scrollIntoView({ behavior: 'smooth' }) } + } + }) - await tick() - - const element = getNodeElement(editor, nodeId) - - if (element) { - element.scrollIntoView({ behavior: 'smooth' }) + const unsubscribeCommentsAddCanceled = documentCommentsAddCanceled.subscribe({ + next: ({ nodeId }) => { + if (editor !== undefined && nodeId != null) { + if (selectNode(editor, nodeId)) { + editor.commands.unsetQMSInlineCommentMark() + } } } }) @@ -104,6 +109,7 @@ onDestroy(() => { unsubscribeHighlightRefresh() unsubscribeNavigateToLocation() + unsubscribeCommentsAddCanceled() }) const handleUpdateTitle = async () => { @@ -137,15 +143,9 @@ return null } - function handleShowDocumentComments (uuid: string) { - if (!uuid) { - return - } - - documentCommentsDisplayRequested({ - element: getNodeElement(editor, uuid), - nodeId: uuid - }) + function handleShowDocumentComments (nodeId: string): void { + const element = getNodeElement(editor, nodeId) + documentCommentsDisplayRequested({ element, nodeId }) } async function createEmbedding (file: File): Promise<{ file: Ref, type: string } | undefined> { @@ -241,32 +241,22 @@ qmsInlineComment: { isHighlightModeOn: () => $canViewDocumentComments || $canAddDocumentComments, getNodeHighlight: handleNodeHighlight, - onNodeSelected: (uuid) => { - if (selectedNodeId !== uuid) { - selectedNodeId = uuid - } - if (isFocused) { - documentCommentsHighlightUpdated(selectedNodeId !== null ? { nodeId: selectedNodeId } : null) - } - }, - onNodeClicked: (uuid) => { - if (selectedNodeId !== uuid) { - selectedNodeId = uuid - } + onNodeClicked: (uuids) => { + // filter out those uuids that are not in comments + uuids = Array.isArray(uuids) ? uuids : [uuids] + uuids = uuids.filter((id) => commentUuids.includes(id)).sort() - if (!$arePopupsOpened && $canViewDocumentComments && selectedNodeId) { + // scroll through the comments as user clicks on the same node + const currIndex = selectedNodeId != null ? uuids.indexOf(selectedNodeId) : -1 + const nextIndex = currIndex === -1 ? 0 : (currIndex + 1) % uuids.length + selectedNodeId = uuids[nextIndex] + + if (!$arePopupsOpened && $canViewDocumentComments && selectedNodeId != null) { handleShowDocumentComments(selectedNodeId) } } } }, - hooks: { - focus: { - onFocus: (focused) => { - isFocused = focused - } - } - }, toc: { onChange: (h) => { headings = h diff --git a/plugins/controlled-documents-resources/src/components/document/popups/AddCommentPopup.svelte b/plugins/controlled-documents-resources/src/components/document/popups/AddCommentPopup.svelte index 1bba64379d..034f56241d 100644 --- a/plugins/controlled-documents-resources/src/components/document/popups/AddCommentPopup.svelte +++ b/plugins/controlled-documents-resources/src/components/document/popups/AddCommentPopup.svelte @@ -11,17 +11,29 @@ const dispatch = createEventDispatcher() - let messageId: Ref = generateId() - async function handleMessage (event: CustomEvent): Promise { + const messageId: Ref = generateId() const comment = await addDocumentCommentFx({ content: event.detail, messageId, nodeId }) - messageId = generateId() dispatch('close', comment) } + + let popup: HTMLDivElement | undefined + + function handleClick (event: MouseEvent): void { + if (event.target instanceof Node) { + if (popup !== undefined && !popup.contains(event.target)) { + event.preventDefault() + event.stopPropagation() + dispatch('close', undefined) + } + } + } -
+ + +
>>( generateActionName('savedAttachmentsUpdated') ) +export const documentCommentsAddCanceled = createEvent<{ + nodeId?: string | null +}>(generateActionName('documentCommentsAddCanceled')) + export const documentCommentsDisplayRequested = createEvent<{ nodeId?: string | null element?: PopupAlignment diff --git a/plugins/controlled-documents-resources/src/stores/editors/document/documentComments.ts b/plugins/controlled-documents-resources/src/stores/editors/document/documentComments.ts index 6365f5dc58..bef590700b 100644 --- a/plugins/controlled-documents-resources/src/stores/editors/document/documentComments.ts +++ b/plugins/controlled-documents-resources/src/stores/editors/document/documentComments.ts @@ -19,9 +19,10 @@ import { type CompAndProps, type PopupAlignment, popupstore, showPopup } from '@ import documents, { type Document, type DocumentComment } from '@hcengineering/controlled-documents' import { isDocumentCommentAttachedTo } from '../../../utils' import { - DocumentCommentPopupCategory, type DocumentCommentsFilter, + DocumentCommentPopupCategory, documentCommentPopupsOpened, + documentCommentsAddCanceled, documentCommentsDisplayRequested, documentCommentsHighlightCleared, documentCommentsHighlightUpdated, @@ -31,7 +32,8 @@ import { documentCommentsSortByChanged, documentCommentsUpdated, controlledDocumentClosed, - savedAttachmentsUpdated + savedAttachmentsUpdated, + controlledDocumentOpened } from './actions' export const $areDocumentCommentPopupsOpened = createStore(false).on( @@ -120,6 +122,7 @@ export const showAddCommentPopupFx = createEffect((payload: { element?: PopupAli payload.element, (result) => { if (result === null || result === undefined) { + documentCommentsAddCanceled({ nodeId: payload.nodeId }) documentCommentsHighlightCleared() } else { documentCommentsDisplayRequested(payload) @@ -187,4 +190,5 @@ export const $savedAttachments = createStore>>([]) .on(savedAttachmentsUpdated, (_, payload) => payload) .reset(controlledDocumentClosed) +forward({ from: controlledDocumentOpened, to: documentCommentsHighlightCleared }) forward({ from: documentCommentsLocationNavigateRequested, to: documentCommentsHighlightUpdated }) diff --git a/plugins/controlled-documents-resources/src/text.ts b/plugins/controlled-documents-resources/src/text.ts index 1b370ce60b..33d340312a 100644 --- a/plugins/controlled-documents-resources/src/text.ts +++ b/plugins/controlled-documents-resources/src/text.ts @@ -24,7 +24,7 @@ import { getCurrentEmployee } from '@hcengineering/contact' import { RequestStatus } from '@hcengineering/request' import { getClient } from '@hcengineering/presentation' import { type ActionContext } from '@hcengineering/text-editor' -import { getNodeElement, selectNode, nodeUuidName } from '@hcengineering/text-editor-resources' +import { getNodeElement, selectNode } from '@hcengineering/text-editor-resources' import { showAddCommentPopupFx } from './stores/editors/document' import { $editorMode } from './stores/editors/document/editor' @@ -127,39 +127,20 @@ async function canAddDocumentComments (doc: ControlledDocument, mode: EditorMode return false } -function setQMSInlineCommentMark (editor: Editor): string | undefined { - if (editor === undefined) { +export async function comment (editor: Editor, event: MouseEvent, ctx: ActionContext): Promise { + const { objectId, objectClass } = ctx + + if (editor === undefined || objectId === undefined || objectClass === undefined) { return } const nodeId = generateId() editor.commands.setQMSInlineCommentMark(nodeId) - return nodeId -} + const element = getNodeElement(editor, nodeId) + await showAddCommentPopupFx({ element, nodeId }) -export async function comment (editor: Editor, event: MouseEvent, ctx: ActionContext): Promise { - const { objectId, objectClass } = ctx - if (objectId === undefined || objectClass === undefined) { - return - } - - let selectedNodeId = editor.extensionStorage[nodeUuidName].activeNodeUuid - - if (selectedNodeId == null) { - selectedNodeId = setQMSInlineCommentMark(editor) - } - - if (selectedNodeId == null) { - return - } - - await showAddCommentPopupFx({ - element: getNodeElement(editor, selectedNodeId), - nodeId: selectedNodeId - }) - - selectNode(editor, selectedNodeId) + selectNode(editor, nodeId) } export async function isCommentVisible (editor: Editor, ctx: ActionContext): Promise { diff --git a/plugins/process-resources/src/components/attributeEditors/FunctionContextPresenter.svelte b/plugins/process-resources/src/components/attributeEditors/FunctionContextPresenter.svelte index 7101d634e6..ae8592de93 100644 --- a/plugins/process-resources/src/components/attributeEditors/FunctionContextPresenter.svelte +++ b/plugins/process-resources/src/components/attributeEditors/FunctionContextPresenter.svelte @@ -22,7 +22,7 @@ export let context: Context const client = getClient() - const func = client.getModel().findObject(contextValue.func) + $: func = client.getModel().findObject(contextValue.func) {#if func !== undefined} diff --git a/plugins/process-resources/src/components/settings/StatesInlineEditor.svelte b/plugins/process-resources/src/components/settings/StatesInlineEditor.svelte index ce66596f54..fcd1884a38 100644 --- a/plugins/process-resources/src/components/settings/StatesInlineEditor.svelte +++ b/plugins/process-resources/src/components/settings/StatesInlineEditor.svelte @@ -17,11 +17,11 @@ import { translate } from '@hcengineering/platform' import { getClient } from '@hcengineering/presentation' import { Process, State } from '@hcengineering/process' + import { makeRank } from '@hcengineering/rank' import { Button, IconAdd, Label } from '@hcengineering/ui' + import { SortableDocList } from '@hcengineering/view-resources' import plugin from '../../plugin' import StateInlineEditor from './StateInlineEditor.svelte' - import { makeRank } from '@hcengineering/rank' - import { SortableList } from '@hcengineering/view-resources' export let process: Process export let states: State[] @@ -49,11 +49,11 @@
- + - + {#if !readonly} - {/each} + + + + +
- + {@const transition = toTransition(value)} - + {#if !readonly}
- {/if} {/if} - {#if isLoading} - - {:else if ($$slots.object ?? presenter) && items} + {#if $$slots.object && items}
{#each items as item, index (item._id)} @@ -225,18 +117,9 @@ }} on:dragend={resetDrag} > - {#if $$slots.object} - - {:else if presenter} - - {/if} +
{/each} - - {#if objectFactory?.component && isCreating} - - (isCreating = false)} /> - {/if} {/if} diff --git a/plugins/view-resources/src/index.ts b/plugins/view-resources/src/index.ts index 7941ffa2c7..b38331d57e 100644 --- a/plugins/view-resources/src/index.ts +++ b/plugins/view-resources/src/index.ts @@ -90,6 +90,7 @@ import YoutubePresenter from './components/linkPresenters/YoutubePresenter.svelt import DividerPresenter from './components/list/DividerPresenter.svelte' import GrowPresenter from './components/list/GrowPresenter.svelte' import ListView from './components/list/ListView.svelte' +import SortableDocList from './components/list/SortableDocList.svelte' import SortableList from './components/list/SortableList.svelte' import SortableListItem from './components/list/SortableListItem.svelte' import TreeElement from './components/navigator/TreeElement.svelte' @@ -236,6 +237,7 @@ export { ObjectIcon, ObjectMention, SortableList, + SortableDocList, SortableListItem, SpaceHeader, SpacePresenter, diff --git a/server-plugins/card-resources/src/index.ts b/server-plugins/card-resources/src/index.ts index ebee2ab0ad..36d5fb046d 100644 --- a/server-plugins/card-resources/src/index.ts +++ b/server-plugins/card-resources/src/index.ts @@ -380,6 +380,37 @@ async function OnCardUpdate (ctx: TxUpdateDoc[], control: TriggerControl): }) } + res.push(...(await updatePeers(control, doc, updateTx))) + return res +} + +async function updatePeers (control: TriggerControl, doc: Card, updateTx: TxUpdateDoc): Promise { + if (updateTx.space === core.space.DerivedTx) return [] + const isDirect = control.hierarchy.isDerived(doc._class, communication.type.Direct) + const isThreadFromDirect = (doc.parentInfo ?? []).some((it) => + control.hierarchy.isDerived(it._class, communication.type.Direct) + ) + + if (!isDirect && !isThreadFromDirect) return [] + + delete updateTx.operations.title + delete updateTx.operations.parentInfo + delete updateTx.operations.parent + delete updateTx.operations.$inc + + const peers = ( + ( + await control.domainRequest(control.ctx, 'communication' as OperationDomain, { + findPeers: { params: { kind: 'card', cardId: doc._id } } + }) + ).value as CardPeer[] + ).flatMap((it) => it.members) + + const res: Tx[] = [] + for (const peer of peers) { + res.push(control.txFactory.createTxUpdateDoc(doc._class, peer.extra.space, peer.cardId, updateTx.operations)) + } + return res } diff --git a/services/github/pod-github/run.sh b/services/github/pod-github/run.sh index 620f7fd770..bf9c19eed8 100755 --- a/services/github/pod-github/run.sh +++ b/services/github/pod-github/run.sh @@ -4,10 +4,8 @@ export CLIENT_SECRET="$POD_GITHUB_CLIENT_SECRET" export PRIVATE_KEY="$POD_GITHUB_PRIVATE_KEY" export SERVER_SECRET=secret export ACCOUNTS_URL=http://localhost:3000 -export COLLABORATOR_URL=http://localhost:3078 -export MINIO_ACCESS_KEY=minioadminchmo -export MINIO_SECRET_KEY=minioadmin -export MINIO_ENDPOINT=localhost -export MONGO_URL=mongodb://localhost:27017 +export COLLABORATOR_URL=ws://huly.local:3078 +export STORAGE_CONFIG="datalake|http://huly.local:4030" +export OTEL_EXPORTER_OTLP_ENDPOINT=http://huly.local:4318/v1/traces rush bundle --to @hcengineering/pod-github node $@ bundle/bundle.js $@ \ No newline at end of file diff --git a/services/github/pod-github/src/platform.ts b/services/github/pod-github/src/platform.ts index 36892bb5d3..7c7d6d4a75 100644 --- a/services/github/pod-github/src/platform.ts +++ b/services/github/pod-github/src/platform.ts @@ -214,6 +214,8 @@ export class PlatformWorker { let oldErrors = '' let sameErrors = 1 + let lastTimeout: any + while (!this.canceled) { let errors: string[] = [] try { @@ -227,6 +229,7 @@ export class PlatformWorker { this.triggerCheckWorkspaces = () => { this.ctx.info('Workspaces check triggered') this.triggerCheckWorkspaces = () => {} + clearTimeout(lastTimeout) resolve() } if (errors.length > 0) { @@ -320,9 +323,10 @@ export class PlatformWorker { try { ;({ client } = await createPlatformClient(ctx, oldWorkspace, 30000)) await this.removeInstallationFromWorkspace(oldWorker, installationId) - await client.close() } catch (err: any) { ctx.error('failed to remove old installation from workspace', { workspace: oldWorkspace, installationId }) + } finally { + await client?.close() } } } @@ -944,26 +948,18 @@ export class PlatformWorker { checkedWorkspaces = new Set() - async checkWorkspaceIsActive ( - workspace: WorkspaceUuid, - workspaceInfo?: WorkspaceInfoWithStatus, - needRecheck = false - ): Promise<{ workspaceInfo: WorkspaceInfoWithStatus | undefined, needRecheck: boolean }> { - if (workspaceInfo === undefined && needRecheck) { - const token = generateToken(systemAccountUuid, workspace, { service: 'github', mode: 'github' }) - workspaceInfo = await getAccountClient(config.AccountsURL, token).getWorkspaceInfo() - } + checkWorkspaceIsActive (workspace: WorkspaceUuid, workspaceInfo: WorkspaceInfoWithStatus): boolean { if (workspaceInfo?.uuid === undefined) { this.ctx.error('No workspace exists for workspaceId', { workspace }) - return { workspaceInfo: undefined, needRecheck: false } + return false } if (workspaceInfo?.isDisabled === true || isDeletingMode(workspaceInfo?.mode)) { this.ctx.warn('Workspace is disabled', { workspace }) - return { workspaceInfo: undefined, needRecheck: false } + return false } if (!isActiveMode(workspaceInfo?.mode)) { this.ctx.warn('Workspace is in maitenance, skipping for now.', { workspace, mode: workspaceInfo?.mode }) - return { workspaceInfo: undefined, needRecheck: true } + return true } const lastVisit = (Date.now() - (workspaceInfo.lastVisit ?? 0)) / (3600 * 24 * 1000) // In days @@ -973,9 +969,38 @@ export class PlatformWorker { this.checkedWorkspaces.add(workspace) this.ctx.warn('Workspace is inactive for too long, skipping for now.', { workspace }) } - return { workspaceInfo: undefined, needRecheck: true } + return true } - return { workspaceInfo, needRecheck: true } + return false + } + + checkReconnect (workspace: WorkspaceUuid, event: ClientConnectEvent, worker: GithubWorker): void { + if (event === ClientConnectEvent.Refresh || event === ClientConnectEvent.Upgraded) { + void this.clients + .get(workspace) + ?.refreshClient(event === ClientConnectEvent.Upgraded) + ?.catch((err) => { + worker.ctx.error('Failed to refresh', { error: err }) + }) + } + + // We need to check if workspace is inactive + const token = generateToken(systemAccountUuid, workspace, { service: 'github', mode: 'github' }) + getAccountClient(config.AccountsURL, token) + .getWorkspaceInfo() + .then((wsInfo) => { + const res = this.checkWorkspaceIsActive(workspace, wsInfo) + if (!res) { + this.ctx.warn('Workspace is inactive, removing from clients list.', { workspace }) + this.clients.delete(workspace) + void worker?.close().catch((err) => { + this.ctx.error('Failed to close workspace', { workspace, error: err }) + }) + } + }) + .catch((err) => { + this.ctx.error('Failed to check workspace is active', { workspace, error: err }) + }) } private async checkWorkspaces (): Promise { @@ -1008,6 +1033,7 @@ export class PlatformWorker { this.ctx.info('connecting to workspace', { workspace: c, time: Date.now() - d.time, version: d.version }) } }, 5000) + try { const token = generateToken(systemAccountUuid, undefined, { service: 'github', mode: 'github' }) const infos = new Map( @@ -1022,82 +1048,57 @@ export class PlatformWorker { toDelete.delete(workspace) continue } + const returnedInfo = infos.get(workspace) + if (returnedInfo === undefined) { + rechecks.push(workspace) + continue + } + const needRecheck = this.checkWorkspaceIsActive(workspace, returnedInfo) + if (needRecheck) { + rechecks.push(workspace) + continue + } await rateLimiter.add(async () => { - const { workspaceInfo, needRecheck } = await this.checkWorkspaceIsActive(workspace, infos.get(workspace)) - if (workspaceInfo === undefined) { - if (needRecheck) { - rechecks.push(workspace) - } - return - } try { - const branding = Object.values(this.brandingMap).find((b) => b.key === workspaceInfo?.branding) ?? null - const workerCtx = this.ctx.newChild('worker', { workspace: workspaceInfo.uuid }, { span: false }) + const branding = Object.values(this.brandingMap).find((b) => b.key === returnedInfo?.branding) ?? null + const workerCtx = this.ctx.newChild('worker', { workspace: returnedInfo.uuid }, { span: false }) - connecting.set(workspaceInfo.uuid, { + connecting.set(returnedInfo.uuid, { time: Date.now(), version: versionToString({ - major: workspaceInfo.versionMajor, - minor: workspaceInfo.versionMinor, - patch: workspaceInfo.versionPatch + major: returnedInfo.versionMajor, + minor: returnedInfo.versionMinor, + patch: returnedInfo.versionPatch }) }) workerCtx.info('************************* Register worker ************************* ', { - workspaceId: workspaceInfo.uuid, - workspaceUrl: workspaceInfo.url, - versionMajor: workspaceInfo.versionMajor, - versionMinor: workspaceInfo.versionMinor, - versionPatch: workspaceInfo.versionPatch, - mode: workspaceInfo.mode, + workspaceId: returnedInfo.uuid, + workspaceUrl: returnedInfo.url, + versionMajor: returnedInfo.versionMajor, + versionMinor: returnedInfo.versionMinor, + versionPatch: returnedInfo.versionPatch, + mode: returnedInfo.mode, index: widx, total: workspaces.length }) - let initialized = false const worker = await GithubWorker.create( this, workerCtx, this.installations, { - dataId: workspaceInfo.dataId, - url: workspaceInfo.url, - uuid: workspaceInfo.uuid + dataId: returnedInfo.dataId, + url: returnedInfo.url, + uuid: returnedInfo.uuid }, branding, this.app, - this.storageAdapter, - (workspace, event) => { - if (event === ClientConnectEvent.Refresh || event === ClientConnectEvent.Upgraded) { - void this.clients - .get(workspace) - ?.refreshClient(event === ClientConnectEvent.Upgraded) - ?.catch((err) => { - workerCtx.error('Failed to refresh', { error: err }) - }) - } - if (initialized) { - // We need to check if workspace is inactive - void this.checkWorkspaceIsActive(workspace, undefined, true) - .then((res) => { - if (res === undefined) { - this.ctx.warn('Workspace is inactive, removing from clients list.', { workspace }) - this.clients.delete(workspace) - void worker?.close().catch((err) => { - this.ctx.error('Failed to close workspace', { workspace, error: err }) - }) - } - }) - .catch((err) => { - this.ctx.error('Failed to check workspace is active', { workspace, error: err }) - }) - } - } + this.storageAdapter ) if (worker !== undefined) { - initialized = true workerCtx.info('************************* Register worker Done ************************* ', { - workspaceId: workspaceInfo.uuid, - workspaceUrl: workspaceInfo.url, + workspaceId: returnedInfo.uuid, + workspaceUrl: returnedInfo.url, index: widx, total: workspaces.length }) @@ -1107,12 +1108,12 @@ export class PlatformWorker { workerCtx.info( '************************* Failed Register worker, timeout or integrations removed *************************', { - workspaceId: workspaceInfo.uuid, - workspaceUrl: workspaceInfo.url, - versionMajor: workspaceInfo.versionMajor, - versionMinor: workspaceInfo.versionMinor, - versionPatch: workspaceInfo.versionPatch, - lastVisit: (Date.now() - (workspaceInfo.lastVisit ?? 0)) / (24 * 60 * 60 * 1000), + workspaceId: returnedInfo.uuid, + workspaceUrl: returnedInfo.url, + versionMajor: returnedInfo.versionMajor, + versionMinor: returnedInfo.versionMinor, + versionPatch: returnedInfo.versionPatch, + lastVisit: (Date.now() - (returnedInfo.lastVisit ?? 0)) / (24 * 60 * 60 * 1000), index: widx, total: workspaces.length } @@ -1124,7 +1125,7 @@ export class PlatformWorker { this.ctx.info("Couldn't create WS worker", { workspace, error: e }) rechecks.push(workspace) } finally { - connecting.delete(workspaceInfo.uuid) + connecting.delete(returnedInfo.uuid) } }) } diff --git a/services/github/pod-github/src/worker.ts b/services/github/pod-github/src/worker.ts index d1ff76ebfb..ca15bb32b4 100644 --- a/services/github/pod-github/src/worker.ts +++ b/services/github/pod-github/src/worker.ts @@ -1766,13 +1766,13 @@ export class GithubWorker implements IntegrationManager { workspace: WorkspaceIds, branding: Branding | null, app: App, - storageAdapter: StorageAdapter, - reconnect: (workspaceId: WorkspaceUuid, event: ClientConnectEvent) => void + storageAdapter: StorageAdapter ): Promise { ctx.info('Connecting to', { workspace }) let client: Client | undefined let endpoint: string | undefined let maitenanceState = false + let worker: GithubWorker | undefined try { ;({ client, endpoint } = await createPlatformClient( ctx, @@ -1784,7 +1784,9 @@ export class GithubWorker implements IntegrationManager { maitenanceState = true throw new Error('Workspace in maintenance') } - reconnect(workspace.uuid, event) + if (worker !== undefined) { + platformWorker.checkReconnect(workspace.uuid, event, worker) + } } )) ctx.info('connected to github', { workspace: workspace.uuid, endpoint }) @@ -1796,7 +1798,7 @@ export class GithubWorker implements IntegrationManager { await GithubWorker.checkIntegrations(client, installations) - const worker = new GithubWorker( + worker = new GithubWorker( ctx, platformWorker.getRateLimiter(endpoint ?? ''), platformWorker, @@ -1808,10 +1810,16 @@ export class GithubWorker implements IntegrationManager { branding ) ctx.info('Init worker', { workspace: workspace.url, workspaceId: workspace.uuid }) - void worker.init() + void worker.init().catch((err) => { + ctx.error('failed to init worker', { error: err, workspace: workspace.uuid }) + void client?.close().catch((err) => { + ctx.error('failed to close client after init error', { error: err, workspace: workspace.uuid }) + }) + }) return worker } catch (err: any) { await client?.close() + void worker?.close() if (maitenanceState) { ctx.info('workspace in maintenance, schedule recheck', { workspace: workspace.uuid, endpoint }) return diff --git a/services/process/src/main.ts b/services/process/src/main.ts index fcadfaa1b7..69cc902d6a 100644 --- a/services/process/src/main.ts +++ b/services/process/src/main.ts @@ -37,6 +37,7 @@ import process, { parseError, processError, ProcessError, + ProcessToDo, State, Step, Transition @@ -169,6 +170,10 @@ async function processExecution (control: ProcessControl, record: ProcessMessage if (transition !== undefined) { await execute(execution, transition, control) } else { + if (record.event === process.trigger.OnToDoRemove) { + const rollbackResult = await checkRollback(control, record, execution) + if (rollbackResult) return + } control.ctx.info('No transition found for event', { event: record.event, execution: execution._id, @@ -177,6 +182,23 @@ async function processExecution (control: ProcessControl, record: ProcessMessage } } +async function checkRollback (control: ProcessControl, record: ProcessMessage, execution: Execution): Promise { + const todo = record.context.todo as ProcessToDo + if (!todo?.withRollback) return false + const rollbackTxes = execution.rollback.pop() ?? [] + for (const tx of rollbackTxes) { + const timeout = setTimeout(() => { + control.ctx.warn('TX HANG', tx) + }, 30000) + await control.client.tx(tx) + clearTimeout(timeout) + } + await control.client.update(execution, { + rollback: execution.rollback + }) + return true +} + async function getTriggerRollback (triggger: TriggerImpl, control: ProcessControl): Promise { if (triggger.rollbackFunc !== undefined) { const rollbackFunc = await getResource(triggger.rollbackFunc) @@ -284,13 +306,12 @@ async function executeTransition (execution: Execution, transition: Transition, execution.rollback.push(rollback) } const executionUpdate = getDiffUpdate(execution, { - rollback: execution.rollback.length > 30 ? execution.rollback.slice(-30) : execution.rollback, - context: execution.context, currentState: state._id, status: isDone ? ExecutionStatus.Done : ExecutionStatus.Active, error: null }) executionUpdate.context = execution.context + executionUpdate.rollback = execution.rollback.length > 30 ? execution.rollback.slice(-30) : execution.rollback res.push(client.txFactory.createTxUpdateDoc(execution._class, execution.space, execution._id, executionUpdate)) res.push( client.txFactory.createTxCreateDoc(process.class.ExecutionLog, execution.space, {