diff --git a/docs/APIDOCUMENTATION.md b/docs/APIDOCUMENTATION.md index e51f98fab9..5beef83710 100644 --- a/docs/APIDOCUMENTATION.md +++ b/docs/APIDOCUMENTATION.md @@ -135,6 +135,8 @@ ##### Notes > This callback is also used during OAuth flows to synchronize member data when the backend returns a different `inbound_member_guid` than the one used to start the flow (e.g., during non-OAuth to OAuth migrations). When this happens, the widget will fetch the new member record and update its internal state to use the new GUID. +> +> `most_recent_job_guid` must change when a new job starts and be `null` before the first one. A `CONNECTED` member whose value is `null`, or unchanged since the widget called `runJob`, keeps the Connecting step waiting. ##### Responses diff --git a/src/utilities/runJobSchedule.js b/src/utilities/runJobSchedule.js index 9b83d5a385..e42ea56116 100644 --- a/src/utilities/runJobSchedule.js +++ b/src/utilities/runJobSchedule.js @@ -30,6 +30,17 @@ const isSafeConflictError = (error) => error?.response?.status === 409 const isConnectedWithoutError = (member) => member?.connection_status === ReadableStatuses.CONNECTED && !member?.error?.error_code +const NOT_STARTED_BY_US = { type: null, previousJobGuid: null } + +// Firefly sets an OAuth member CONNECTED on the redirect before any job exists, and over +// websockets that update can arrive after we started ours (CT-2332). It names the job the +// member had before runJob: none for a first job, the previous job for a returning member. +const isPreJobUpdate = (member, started) => { + const guid = member?.most_recent_job_guid + + return !guid || guid === started.previousJobGuid +} + /** * Work out which job just finished, in order of trust: * - the job we loaded fresh off the polled member @@ -102,17 +113,26 @@ export const runJobSchedule$ = ({ * else scheduled after it, otherwise we keep polling until the member is idle * so the next job can be started. */ - const observeRunningJob = (memberGuid, currentSchedule, startedType) => + const observeRunningJob = (memberGuid, currentSchedule, started) => pollMember(memberGuid).pipe( + // onPoll runs before the gate on purpose: it is where the Connecting timeout lives. tap(onPoll), - filter((pollingState) => pollingState.pollingIsDone), + // Error and MFA states route on the member alone; only CONNECTED needs a real finished job. + filter((pollingState) => { + const polledMember = pollingState.currentResponse?.member + + return ( + pollingState.pollingIsDone && + !(isConnectedWithoutError(polledMember) && isPreJobUpdate(polledMember, started)) + ) + }), take(1), map((pollingState) => pollingState.currentResponse), mergeMap((polledResponse) => loadJob(polledResponse.member).pipe( map((job) => ({ member: polledResponse.member, - job: resolveFinishedJob(job, polledResponse.job, startedType), + job: resolveFinishedJob(job, polledResponse.job, started.type), })), ), ), @@ -135,8 +155,8 @@ export const runJobSchedule$ = ({ }), ) - const observeThenContinue = (memberGuid, currentSchedule, iteration, startedType) => - observeRunningJob(memberGuid, currentSchedule, startedType).pipe( + const observeThenContinue = (memberGuid, currentSchedule, iteration, started) => + observeRunningJob(memberGuid, currentSchedule, started).pipe( mergeMap(({ member: observedMember, job }) => { const emitted = of({ member: observedMember, job }) @@ -157,22 +177,27 @@ export const runJobSchedule$ = ({ } if (currentMember.is_being_aggregated !== false) { - return observeThenContinue(currentMember.guid, currentSchedule, iteration, null) + return observeThenContinue( + currentMember.guid, + currentSchedule, + iteration, + NOT_STARTED_BY_US, + ) } const activeJob = JobSchedule.getActiveJob(currentSchedule) return defer(() => api.runJob(activeJob.type, currentMember.guid, config, true)).pipe( - map(() => activeJob.type), + map(() => ({ type: activeJob.type, previousJobGuid: currentMember.most_recent_job_guid })), catchError((error) => { // 409 is usually the job Firefly created on the OAuth redirect. // It gets observed and reconciled like any other running job. - if (isSafeConflictError(error)) return of(null) + if (isSafeConflictError(error)) return of(NOT_STARTED_BY_US) return throwError(() => error) }), - mergeMap((startedType) => - observeThenContinue(currentMember.guid, currentSchedule, iteration, startedType), + mergeMap((started) => + observeThenContinue(currentMember.guid, currentSchedule, iteration, started), ), ) }) diff --git a/src/utilities/runJobSchedule.test.tsx b/src/utilities/runJobSchedule.test.tsx new file mode 100644 index 0000000000..6b4ee75856 --- /dev/null +++ b/src/utilities/runJobSchedule.test.tsx @@ -0,0 +1,147 @@ +import { waitFor } from 'src/utilities/testingLibrary' +import { POST_MESSAGES } from 'src/const/postMessages' +import { ReadableStatuses } from 'src/const/Statuses' +import { JOB_TYPES } from 'src/const/consts' +import { STEPS, VERIFY_MODE } from 'src/const/Connect' +import { ACTIONABLE_ERROR_CODES } from 'src/views/actionableError/consts' +import { + createFakeBackend, + createFakeBrokaw, + expectMemberConnected, + HttpError, + Member, + REDIRECT_JOB_GUID, + renderConnecting, + staleOAuthMember, +} from 'src/utilities/test/connectingOAuthHarness' + +// fadeOut (Velocity) never resolves in jsdom; Connecting's error path dispatches inside its .then. +vi.mock('src/utilities/Animation', () => ({ fadeOut: vi.fn(() => Promise.resolve()) })) + +/** + * runJobSchedule$ drives the Connecting step's job schedule. These tests run it through the + * real (real store, hook and transport); only the API and brokaw are faked. + * + * CT-2332: firefly sets an OAuth member CONNECTED on the redirect before any job exists. + * Over websockets, a copy of that member update can reach the widget *after* it has started + * its own job. It looks finished (CONNECTED, not aggregating) but names no job, or the job the + * member had before. The widget must not mistake it for its job finishing. + */ + +const OUR_JOB_GUID = `JOB-${JOB_TYPES.VERIFICATION}` + +// The member as firefly leaves it on the OAuth redirect: CONNECTED, idle, no job yet. +const connectedWithNoJob: Member = { + ...staleOAuthMember, + connection_status: ReadableStatuses.CONNECTED, +} + +const runningJob = (jobGuid: string): Member => ({ + ...connectedWithNoJob, + is_being_aggregated: true, + most_recent_job_guid: jobGuid, +}) + +const finishedJob = (jobGuid: string): Member => ({ + ...connectedWithNoJob, + most_recent_job_guid: jobGuid, +}) + +const impededWithNoEligibleAccounts = (jobGuid: string): Member => ({ + ...staleOAuthMember, + connection_status: ReadableStatuses.IMPEDED, + most_recent_job_guid: jobGuid, + error: { error_code: ACTIONABLE_ERROR_CODES.NO_ELIGIBLE_ACCOUNTS }, +}) + +const expectNoEligibleAccountsScreen = async (widget: ReturnType) => { + // Wait for Connecting to route anywhere, then check where. Before the fix it routed to + // CONNECTED as soon as the stale update arrived. + await waitFor(() => expect(widget.currentStep()).toBeDefined(), { timeout: 4000 }) + expect(widget.currentStep()).toBe(STEPS.ACTIONABLE_ERROR) + expect(widget.onPostMessage).not.toHaveBeenCalledWith( + POST_MESSAGES.MEMBER_CONNECTED, + expect.anything(), + ) +} + +describe('runJobSchedule$ through over websockets', () => { + afterEach(() => { + vi.restoreAllMocks() + }) + + it('ignores the late CONNECTED update with no job and shows the real outcome of the job it started', async () => { + // Given a first-time OAuth member: no job has ever run on it. + const backend = createFakeBackend() + const brokaw = createFakeBrokaw() + const widget = renderConnecting( + backend, + { mode: VERIFY_MODE }, + { webSocket: brokaw.connection }, + ) + + // When the widget starts its verification job... + await widget.runJobCalled() + // ...and firefly's pre-job update arrives late, looking finished but naming no job... + await brokaw.memberUpdated(connectedWithNoJob) + // ...then the widget's job actually finishes with no eligible accounts. + await brokaw.memberUpdated(impededWithNoEligibleAccounts(OUR_JOB_GUID)) + + // Then the widget shows the error, never a success. + await expectNoEligibleAccountsScreen(widget) + }) + + it('ignores the late CONNECTED update that still names a returning member’s previous job', async () => { + // Given a returning member whose previous job was also a verification. Attributing that + // old job by type would wrongly complete the schedule. + const PREVIOUS_JOB_GUID = 'JOB-old' + const returningMember = finishedJob(PREVIOUS_JOB_GUID) + const backend = createFakeBackend({ member: returningMember }) + backend.jobs[PREVIOUS_JOB_GUID] = { guid: PREVIOUS_JOB_GUID, job_type: JOB_TYPES.VERIFICATION } + const brokaw = createFakeBrokaw() + const widget = renderConnecting( + backend, + { mode: VERIFY_MODE }, + { webSocket: brokaw.connection, member: returningMember }, + ) + + // When the widget starts a new verification job... + await widget.runJobCalled() + // ...and firefly's pre-job update arrives late, still naming the previous job... + await brokaw.memberUpdated(returningMember) + // ...then the new job finishes with no eligible accounts. + await brokaw.memberUpdated(impededWithNoEligibleAccounts(OUR_JOB_GUID)) + + // Then the widget shows the error, never a success. + await expectNoEligibleAccountsScreen(widget) + }) + + it('observes the job firefly assigned when its own runJob is rejected with a 409', async () => { + // Given firefly already started the verification job on the redirect + // (disable_background_agg clients), so the widget's own runJob is a duplicate. + const backend = createFakeBackend() + backend.runJob.mockImplementationOnce(async () => { + backend.startJob(REDIRECT_JOB_GUID, JOB_TYPES.VERIFICATION) + throw new HttpError(409) + }) + const brokaw = createFakeBrokaw() + const widget = renderConnecting( + backend, + { mode: VERIFY_MODE }, + { webSocket: brokaw.connection }, + ) + + // When the widget's runJob is rejected... + await widget.runJobCalled() + // ...the late pre-job update is still ignored... + await brokaw.memberUpdated(connectedWithNoJob) + // ...and firefly's job is seen running, then finishing. + await brokaw.memberUpdated(runningJob(REDIRECT_JOB_GUID)) + await brokaw.memberUpdated(finishedJob(REDIRECT_JOB_GUID)) + + // Then the widget completes against firefly's job without starting another. + await expectMemberConnected(widget.onPostMessage) + expect(backend.runJob).toHaveBeenCalledTimes(1) + expect(backend.loadJob).toHaveBeenCalledWith(REDIRECT_JOB_GUID) + }) +}) diff --git a/src/utilities/test/connectingOAuthHarness.tsx b/src/utilities/test/connectingOAuthHarness.tsx new file mode 100644 index 0000000000..c1d15ac41e --- /dev/null +++ b/src/utilities/test/connectingOAuthHarness.tsx @@ -0,0 +1,236 @@ +import React from 'react' +import { Subject } from 'rxjs' +import { render, waitFor } from 'src/utilities/testingLibrary' +import { Connecting } from 'src/views/connecting/Connecting' +import { PostMessageContext } from 'src/ConnectWidget' +import { ApiContextTypes } from 'src/context/ApiContext' +import { WebSocketConnection, WebSocketProvider } from 'src/context/WebSocketContext' +import { POST_MESSAGES } from 'src/const/postMessages' +import { ReadableStatuses } from 'src/const/Statuses' + +/** + * Drives the real after OAuth against fakes for the two things a test cannot + * use for real: the backend (firefly/persona) and brokaw (websockets). + */ + +const MEMBER_GUID = 'MBR-oauth' +const USER_GUID = 'USR-1' +export const REDIRECT_JOB_GUID = 'JOB-redirect' + +export type Member = { + guid: string + user_guid: string + connection_status: number + is_being_aggregated: boolean + most_recent_job_guid: string | null + is_oauth: boolean + error?: { error_code: number } | null +} + +type Job = { guid: string; job_type: number; async_account_data_ready?: boolean } + +export class HttpError extends Error { + response: { status: number } + + constructor(status: number, message = 'Request failed') { + super(message) + this.name = 'HttpError' + this.response = { status } + } +} + +// The member the widget holds when it lands on Connecting: created before the user left for +// the institution, so PENDING and without a job. +export const staleOAuthMember: Member = { + guid: MEMBER_GUID, + user_guid: USER_GUID, + connection_status: ReadableStatuses.PENDING, + is_being_aggregated: false, + most_recent_job_guid: null, + is_oauth: true, +} + +/** + * In-memory firefly/persona. `startJob` puts the member into aggregation for that job, + * `loadMemberByGuid` lets it go idle after `pollsUntilDone` polls, and `runJob` answers 409 + * while a job is already running, as firefly does. + */ +export const createFakeBackend = ({ + pollsUntilDone = 2, + earlyDataRelease = false, + member = staleOAuthMember, +} = {}) => { + const backend = { + member: { ...member } as Member, + jobs: {} as Record, + pollsWhileRunning: 0, + + startJob(guid: string, jobType: number) { + // With early data release the job reports its data as ready while it is + // still running, which makes member polling stop before the job finishes. + backend.jobs[guid] = { guid, job_type: jobType, async_account_data_ready: earlyDataRelease } + backend.member = { + ...backend.member, + connection_status: ReadableStatuses.CONNECTED, + is_being_aggregated: true, + most_recent_job_guid: guid, + } + }, + + loadMemberByGuid: vi.fn(async (): Promise => { + if (backend.member.is_being_aggregated) { + backend.pollsWhileRunning += 1 + + if (backend.pollsWhileRunning >= pollsUntilDone) { + backend.pollsWhileRunning = 0 + backend.member = { ...backend.member, is_being_aggregated: false } + } + } + + return backend.member + }), + + loadJob: vi.fn(async (guid: string): Promise => { + const job = backend.jobs[guid] + + if (!job) { + throw new HttpError(404) + } + + return job + }), + + runJob: vi.fn(async (jobType: number): Promise> => { + if (backend.member.is_being_aggregated) { + throw new HttpError(409) + } + + backend.startJob(`JOB-${jobType}`, jobType) + + return {} + }), + } + + return backend +} + +// Yields one macrotask so the widget's pending promises and subscriptions settle. +const settle = () => new Promise((resolve) => setTimeout(resolve, 0)) + +/** + * Stands in for brokaw. `memberUpdated(member)` delivers a `members/updated` frame and waits + * for the widget to process it. Like brokaw, nothing is replayed to late subscribers, so send + * frames only once the widget is observing. + */ +export const createFakeBrokaw = () => { + const frames$ = new Subject<{ event: string; payload: Member }>() + + const connection: WebSocketConnection = { + isConnected: () => true, + webSocketMessages$: frames$.asObservable(), + } + + const memberUpdated = async (member: Member) => { + frames$.next({ event: 'members/updated', payload: member }) + await settle() + } + + return { connection, memberUpdated } +} + +// Connecting throws `connectingError` during render so the host's error boundary can take +// over. Tests need a boundary of their own to observe that. +class TestErrorBoundary extends React.Component< + { onError: (error: Error) => void; children: React.ReactNode }, + { hasError: boolean } +> { + state = { hasError: false } + + componentDidCatch(error: Error) { + this.props.onError(error) + } + + static getDerivedStateFromError() { + return { hasError: true } + } + + render() { + return this.state.hasError ?
: this.props.children + } +} + +export const renderConnecting = ( + backend: ReturnType, + connectConfig: Record, + { + webSocket, + member = staleOAuthMember, + }: { webSocket?: WebSocketConnection; member?: Member } = {}, +) => { + const onPostMessage = vi.fn() + const onError = vi.fn() + + // The shared render helper hard-codes a no-op onPostMessage and has no websocket context, + // so those two are provided here. + const connecting = ( + + + + ) + + const { store } = render( + + {webSocket ? ( + {connecting} + ) : ( + connecting + )} + , + { + apiValue: { + loadMemberByGuid: backend.loadMemberByGuid, + loadJob: backend.loadJob, + runJob: backend.runJob, + } as unknown as ApiContextTypes, + preloadedState: { + connect: { + currentMemberGuid: MEMBER_GUID, + members: [member], + jobSchedule: { isInitialized: false, jobs: [] }, + location: [], + selectedInstitution: {}, + }, + experimentalFeatures: { + // With websockets on, frames drive the observation and polling is effectively off. + memberPollingMilliseconds: webSocket ? 60_000 : 10, + optOutOfEarlyUserRelease: false, + unavailableInstitutions: [], + useWebSockets: !!webSocket, + }, + }, + }, + ) + + // Resolves once the widget has asked the backend to run a job and is observing the result. + const runJobCalled = async () => { + await waitFor(() => expect(backend.runJob).toHaveBeenCalled()) + await settle() + } + + const currentStep = () => { + const { location } = store.getState().connect + return location[location.length - 1]?.step + } + + return { onPostMessage, onError, runJobCalled, currentStep } +} + +export const expectMemberConnected = (onPostMessage: ReturnType) => + waitFor( + () => + expect(onPostMessage).toHaveBeenCalledWith(POST_MESSAGES.MEMBER_CONNECTED, { + user_guid: USER_GUID, + member_guid: MEMBER_GUID, + }), + { timeout: 5000 }, + ) diff --git a/src/views/connecting/__tests__/ConnectingOAuthJobs-test.tsx b/src/views/connecting/__tests__/ConnectingOAuthJobs-test.tsx index 7589204f9b..85b1f6f359 100644 --- a/src/views/connecting/__tests__/ConnectingOAuthJobs-test.tsx +++ b/src/views/connecting/__tests__/ConnectingOAuthJobs-test.tsx @@ -1,13 +1,17 @@ -import React from 'react' -import { createTestReduxStore, render, waitFor } from 'src/utilities/testingLibrary' -import { Connecting } from 'src/views/connecting/Connecting' -import { PostMessageContext } from 'src/ConnectWidget' -import { ApiContextTypes, ApiProvider } from 'src/context/ApiContext' +import { waitFor } from 'src/utilities/testingLibrary' import { POST_MESSAGES } from 'src/const/postMessages' import { ReadableStatuses } from 'src/const/Statuses' import { JOB_TYPES } from 'src/const/consts' import { VERIFY_MODE } from 'src/const/Connect' import { EXTRA_ITERATIONS_ALLOWED } from 'src/utilities/runJobSchedule' +import { + createFakeBackend, + expectMemberConnected, + HttpError, + REDIRECT_JOB_GUID, + renderConnecting, + staleOAuthMember, +} from 'src/utilities/test/connectingOAuthHarness' /** * CT-2495: after OAuth the widget lands on Connecting holding the member it @@ -18,174 +22,6 @@ import { EXTRA_ITERATIONS_ALLOWED } from 'src/utilities/runJobSchedule' * that used to leave the widget on this screen forever. */ -const MEMBER_GUID = 'MBR-oauth' -const USER_GUID = 'USR-1' -const REDIRECT_JOB_GUID = 'JOB-redirect' - -type Member = { - guid: string - user_guid: string - connection_status: number - is_being_aggregated: boolean - most_recent_job_guid: string | null - is_oauth: boolean -} - -type Job = { guid: string; job_type: number; async_account_data_ready?: boolean } - -class HttpError extends Error { - response: { status: number } - - constructor(status: number, message = 'Request failed') { - super(message) - this.name = 'HttpError' - this.response = { status } - } -} - -const staleOAuthMember: Member = { - guid: MEMBER_GUID, - user_guid: USER_GUID, - connection_status: ReadableStatuses.PENDING, - is_being_aggregated: false, - most_recent_job_guid: null, - is_oauth: true, -} - -const connectedMemberRunning = (jobGuid: string): Member => ({ - ...staleOAuthMember, - connection_status: ReadableStatuses.CONNECTED, - is_being_aggregated: true, - most_recent_job_guid: jobGuid, -}) - -const createStore = () => - createTestReduxStore({ - connect: { - currentMemberGuid: MEMBER_GUID, - members: [staleOAuthMember], - jobSchedule: { isInitialized: false, jobs: [] }, - location: [], - selectedInstitution: {}, - }, - experimentalFeatures: { - memberPollingMilliseconds: 10, - optOutOfEarlyUserRelease: false, - unavailableInstitutions: [], - useWebSockets: false, - }, - }) - -const createFakeBackend = ({ pollsUntilDone = 2, earlyDataRelease = false } = {}) => { - const backend = { - member: { ...staleOAuthMember } as Member, - jobs: {} as Record, - pollsWhileRunning: 0, - - startJob(guid: string, jobType: number) { - // With early data release the job reports its data as ready while it is - // still running, which makes member polling stop before the job finishes. - backend.jobs[guid] = { guid, job_type: jobType, async_account_data_ready: earlyDataRelease } - backend.member = connectedMemberRunning(guid) - }, - - loadMemberByGuid: vi.fn(async (): Promise => { - if (backend.member.is_being_aggregated) { - backend.pollsWhileRunning += 1 - - if (backend.pollsWhileRunning >= pollsUntilDone) { - backend.pollsWhileRunning = 0 - backend.member = { ...backend.member, is_being_aggregated: false } - } - } - - return backend.member - }), - - loadJob: vi.fn(async (guid: string): Promise => { - const job = backend.jobs[guid] - - if (!job) { - throw new HttpError(404) - } - - return job - }), - - runJob: vi.fn(async (jobType: number): Promise> => { - if (backend.member.is_being_aggregated) { - // Firefly returns a 409 when the member already has a running job. - throw new HttpError(409) - } - - backend.startJob(`JOB-${jobType}`, jobType) - - return {} - }), - } - - return backend -} - -/** - * Connecting throws `connectingError` during render so the host's error - * boundary can take over. Tests need a boundary of their own to observe that. - */ -class TestErrorBoundary extends React.Component< - { onError: (error: Error) => void; children: React.ReactNode }, - { hasError: boolean } -> { - state = { hasError: false } - - componentDidCatch(error: Error) { - this.props.onError(error) - } - - static getDerivedStateFromError() { - return { hasError: true } - } - - render() { - return this.state.hasError ?
: this.props.children - } -} - -const renderConnecting = ( - backend: ReturnType, - connectConfig: Record, -) => { - const onPostMessage = vi.fn() - const onError = vi.fn() - const api = { - loadMemberByGuid: backend.loadMemberByGuid, - loadJob: backend.loadJob, - runJob: backend.runJob, - } as unknown as ApiContextTypes - - render( - - - - - - - , - { store: createStore() }, - ) - - return { onPostMessage, onError } -} - -const expectMemberConnected = (onPostMessage: ReturnType) => - waitFor( - () => - expect(onPostMessage).toHaveBeenCalledWith(POST_MESSAGES.MEMBER_CONNECTED, { - user_guid: USER_GUID, - member_guid: MEMBER_GUID, - }), - { timeout: 5000 }, - ) - describe(' after OAuth', () => { afterEach(() => { vi.restoreAllMocks()