verdnatura-chat/ios/Pods/Flipper-RSocket/yarpl/flowable/FlowableOperator.h

1011 lines
30 KiB
C
Raw Normal View History

Merge beta into master (#2143) * [FIX] Messages being sent but showing as temp status (#1469) * [FIX] Missing messages after reconnect (#1470) * [FIX] Few fixes on themes (#1477) * [I18N] Missing German translations (#1465) * Missing German translation * adding a missing space behind colon * added a missing space after colon * and another attempt to finally fix this – got confused by all the branches * some smaller fixes for the translation * better wording * fixed another typo * [FIX] Crash while displaying the attached image with http on file name (#1401) * [IMPROVEMENT] Tap app and server version to copy to clipboard (#1425) * [NEW] Reply notification (#1448) * [FIX] Incorrect background color login on iPad (#1480) * [FIX] Prevent multiple tap on send (Share Extension) (#1481) * [NEW] Image Viewer (#1479) * [DOCS] Update Readme (#1485) * [FIX] Jitsi with Hermes Enabled (#1523) * [FIX] Draft messages not working with themed Messagebox (#1525) * [FIX] Go to direct message from members list (#1519) * [FIX] Make SAML wait for idp token instead of creating it on client (#1527) * [FIX] Server Test Push Notification (#1508) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [CHORE] Update to new server response (#1509) * [FIX] Insert messages with blank users (#1529) * Bump version to 4.2.1 (#1530) * [FIX] Error when normalizing empty messages (#1532) * [REGRESSION] CAS (#1570) * Bump version to 4.2.2 (#1571) * [FIX] Add username block condition to prevent error (#1585) * Bump version to 4.2.3 * Bump version to 4.2.4 * Bump version to 4.3.0 (#1630) * [FIX] Channels doesn't load (#1586) * [FIX] Channels doesn't load * [FIX] Update roomsUpdatedAt when subscriptions.length is 0 * [FIX] Remove unnecessary changes * [FIX] Improve the code Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Make SAML to work on Rocket.Chat < 2.3.0 (#1629) * [NEW] Invite links (#1534) * [FIX] Set the http-agent to the form that Rocket.Chat requires for logging (#1482) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] "Following thread" and "Unfollowed Thread" is hardcoded and not translated (#1625) * [FIX] Disable reset button if form didn't changed (#1569) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Header title of RoomInfoView (#1553) * [I18N] Gallery Permissions DE (#1542) * [FIX] Not allow to send messages to archived room (#1623) * [FIX] Profile fields automatically reset (#1502) * [FIX] Show attachment on ThreadMessagesView (#1493) * [NEW] Wordpress auth (#1633) * [CHORE] Add Start Packager script (#1639) * [CHORE] Update RN to 0.61.5 (#1638) * [CHORE] Update RN to 0.61.5 * [CHORE] Update react-native patch Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * Bump version to 4.3.1 (#1641) * [FIX] Change force logout rule (#1640) * Bump version to 4.4.0 (#1643) * [IMPROVEMENT] Use MessagingStyle on Android Notification (#1575) * [NEW] Request review (#1627) * [NEW] Pull to refresh RoomView (#1657) * [FIX] Unsubscribe from room (#1655) * [FIX] Server with subdirs (#1646) * [NEW] Clear cache (#1660) * [IMPROVEMENT] Memoize and batch subscriptions updates (#1642) * [FIX] Disallow empty sharing (#1664) * [REGRESSION] Use HTTPS links for sharing and markets protocol for review (#1663) * [FIX] In some cases, share extension doesn't load images (#1649) * [i18n] DE translations for new invite function and some minor fixes (#1631) * [FIX] Remove duplicate jetify step (#1628) minor: also remove 'cd' calls Co-authored-by: Diego Mello <diegolmello@gmail.com> * [REGRESSION] Read messages (#1666) * [i18n] German translations missing (#1670) * [FIX] Notifications crash on older Android Versions (#1672) * [i18n] Added Dutch translation (#1676) * [NEW] Omnichannel Beta (#1674) * [NEW] Confirm logout/clear cache (#1688) * [I18N] Add es-ES language (#1495) * [NEW] UiKit Beta (#1497) * [IMPROVEMENT] Use reselect (#1696) * [FIX] Notification in Android API level less than 24 (#1692) * [IMPROVEMENT] Send tmid on slash commands and media (#1698) * [FIX] Unhandled action on UIKit (#1703) * [NEW] Pull to refresh RoomsList (#1701) * [IMPROVEMENT] Reset app when language is changed (#1702) * [FIX] Small fixes on UIKit (#1709) * [FIX] Spotlight (#1719) * [CHORE] Update react-native-image-crop-picker (#1712) * [FIX] Messages Overlapping (Android) and MessageBox Scroll (iOS) (#1720) * [REGRESSION] Remove @ and # from mention (#1721) * [NEW] Direct message from user info (#1516) * [FIX] Delete slash commands (#1723) * [IMPROVEMENT] Hold URL to copy (#1684) * [FIX] Different sourcemaps generation for Hermes (#1724) * [FIX] Different sourcemaps generation for Hermes * Upload sourcemaps after build * [REVERT] Show emoji keyboard on Android (#1738) * [FIX] Stop logging react-native-image-crop-picker (#1745) * [FIX] Prevent toast ref error (#1744) * [FIX] Prevent reaction map error (#1743) * [FIX] Add missing calls to user info (#1741) * [FIX] Catch room unsubscribe error (#1739) * [i18n] Missing German keys (#1735) * [FIX] Missing i18n on MessagesView title (#1733) * [FIX] UIKit Modal: Weird behavior on Android Tablet (#1742) * [i18n] Missing key on German (#1747) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [i18n] Add Italian (#1736) * [CHORE] Memory leaks investigation (#1675) * [IMPROVEMENT] Alert verify email when enabled (#1725) * [NEW] Jitsi JWT added to URL (#1746) * [FIX] UIKit submit when connection lost (#1748) * Bump version to 4.5.0 (#1761) * [NEW] Default browser (#1752) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] HTTP Basic Auth (#1753) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Honor profile fields edit settings (#1687) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Room announcements (#1726) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Honor Register/Login settings (#1727) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Make links clickable on Room Info (#1730) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [NEW] Hide system messages (#1755) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Honor "Message_AudioRecorderEnabled" (#1764) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [i18n] Missing de keys (#1765) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Redirect user to SetUsernameView (#1728) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Join Room (#1769) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Accept all media types using * (#1770) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Use RealName when necessary (#1758) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Markdown Line Break (#1783) * [IMPROVEMENT] Remove useMarkdown (#1774) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Open browser rather than webview on Create Workspace (#1788) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Markdown perf (#1796) * [FIX] Stop video when modal is closed (#1787) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Hide reply notification action when there are missing data (#1771) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [i18n] Added Japanese translation (#1781) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Reset password error message (#1772) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Close tablet modal (#1773) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Setting not present (#1775) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Thread header (#1776) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Keyboard tracking loses input ref (#1784) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [NEW] Mark message as unread (#1785) Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [IMPROVEMENT] Log server version (#1786) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Add loading message on long running tasks (#1798) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [CHORE] Switch Apple account on Fastlane (#1810) * [FIX] Watermelon throwing "Cannot update a record with pending updates" (#1754) * [FIX] Detox tests (#1790) * [CHORE] Use markdown preview on RoomView Header (#1807) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] LoginSignup blink services (#1809) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Request user presence on demand (#1813) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Remove all invited users when create a channel (#1814) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Pop from room which you have been removed (#1819) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Room Info styles (#1820) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [i18n] Add missing German keys (#1800) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Empty mentions for @all and @here when real name is enabled (#1822) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [TESTS] Markdown added to Storybook (#1812) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [REGRESSION] Room View header title (#1827) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Storybook snapshots (#1831) Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [FIX] Mentions (#1829) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Thread message not found (#1830) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Separate delete and remove channel (#1832) * Rename to delete room * Separate delete and remove channel * handleRemoved -> handleRoomRemoved * [FIX] Navigate to RoomsList & Handle tablet case Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [NEW] Filter system messages per room (#1815) Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] e2e tests (#1838) * [FIX] Consecutive clear cache calls freezing app (#1851) * Bump version to 4.5.1 (#1853) * [FIX][iOS] Ignore silent mode on audio player (#1862) * [IMPROVEMENT] Create App Group property on Info.plist (#1858) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Make username clickable on message (#1618) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Show proper error message on profile (#1768) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Show toast when a message is starred/unstarred (#1616) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Incorrect size params to avatar endpoint (#1875) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Remove unrecognized emoji flags on android (#1887) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Remove react-native global installs (#1886) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Emojis transparent on android (#1881) Co-authored-by: Diego Mello <diegolmello@gmail.com> * Bump acorn from 5.7.3 to 5.7.4 (#1876) Bumps [acorn](https://github.com/acornjs/acorn) from 5.7.3 to 5.7.4. - [Release notes](https://github.com/acornjs/acorn/releases) - [Commits](https://github.com/acornjs/acorn/compare/5.7.3...5.7.4) Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * Bump version to 4.6.0 (#1911) * [FIX] Encode Image URI (#1909) * [FIX] Encode Image URI * [FIX] Check if Image is Valid Co-authored-by: Diego Mello <diegolmello@gmail.com> * [NEW] Adaptive Icons (#1904) * Remove unnecessary stuff from debug build * Adaptive icon for experimental app * [FIX] Stop showing message on leave channel (#1896) * [FIX] Leave room don't show 'was removed' message * [FIX] Remove duplicated code Co-authored-by: Diego Mello <diegolmello@gmail.com> * [i18n] Added missing German translations(#1900) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Linkedin OAuth login (#1913) * [CHORE] Fix typo in CreateChannel View (#1930) * [FIX] Respect protocol in HTTP Auth IPs (#1933) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Use new LinkedIn OAuth url (#1935) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [CHORE] Use storyboard on splash screen (#1939) * Update react-native-bootsplash * iOS * Fix android * [FIX] Check if avatar exists before create Icon (#1927) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Ignore self typing event (#1950) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Change default directory listing to Users (#1948) * fix: change default directory listing to Users * follow server settings * Fix state to props Co-authored-by: Diego Mello <diegolmello@gmail.com> * [NEW] Onboarding layout (#1954) * Onboarding texts * OnboardingView * FormContainer * Minor fixes * NewServerView * Remove code * Refactor * WorkspaceView * Stash * Login with email working * Login with * Join open * Revert "Login with" This reverts commit d05dc507d2e9a2db76d433b9b1f62192eba35dbd. * Fix create account styles * Register * Refactor * LoginServices component * Refactor * Multiple servers * Remove native images * Refactor styles * Fix testid * Fix add server on tablet * i18n * Fix close modal * Fix TOTP * [FIX] Registration disabled * [FIX] Login Services separator * Fix logos * Fix AppVersion name * I18n * Minor fixes * [FIX] Custom Fields Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [NEW] Create discussions (#1942) * [WIP][NEW] Create Discussion * [FIX] Clear multiselect & Translations * [NEW] Create Discussion at MessageActions * [NEW] Disabled Multiselect * [FIX] Initial channel * [NEW] Create discussion on MessageBox Actions * [FIX] Crashing on edit name * [IMPROVEMENT] New message layout * [CHORE] Update README * [NEW] Avatars on MultiSelect * [FIX] Select Users * [FIX] Add redirect and Handle tablet * [IMPROVEMENT] Split CreateDiscussionView * [FIX] Create a discussion inner discussion * [FIX] Create a discussion * [I18N] Add pt-br * Change icons * [FIX] Nav to discussion & header title * Fix header Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Load messages (#1910) * Create updateLastOpen param on readMessages * Remove InteractionManager from load messages * [NEW] Custom Status (#1811) * [NEW] Custom Status * [FIX] Subscribe to changes * [FIX] Improve code using Banner component * [IMPROVEMENT] Toggle modal * [NEW] Edit custom status from Sidebar * [FIX] Modal when tablet * [FIX] Styles * [FIX] Switch to react-native-promp-android * [FIX] Custom Status UI * [TESTS] E2E Custom Status * Fix banner * Fix banner * Fix subtitle * status text * Fix topic header * Fix RoomActionsView topic * Fix header alignment on Android * [FIX] RoomInfo crashes when without statusText * [FIX] Use users.setStatus * [FIX] Remove customStatus of ProfileView * [FIX] Room View Thread Header Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] UI issues of Create Discussion View (#1965) * [NEW] Direct Message between multiple users (#1958) * [WIP] DM between multiple users * [WIP][NEW] Create new DM between multiple users * [IMPROVEMENT] Improve createChannel Sagas * [IMPROVEMENT] Selected Users view * [IMPROVEMENT] Room Actions of Group DM * [NEW] Create new DM between multiple users * [NEW] Group DM avatar * [FIX] Directory border * [IMPROVEMENT] Use isGroupChat * [CHORE] Remove legacy getRoomMemberId * [NEW] RoomTypeIcon * [FIX] No use legacy method on RoomInfoView * [FIX] Blink header when create new DM * [FIX] Only show create direct message option when allowed * [FIX] RoomInfoView * pt-BR * Few fixes * Create button name * Show create button only after a user is selected * Fix max users issues Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Add server and hide login (#1968) * Navigate to new server workspace from ServerDropdown if there's no token * Hide login button based on login services and Accounts_ShowFormLogin setting * [FIX] Lint Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [FIX] MultiSelect Keyboard behavior (Android) (#1969) * fixed-modal-position * made-changes Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [FIX] Bottom border style on DirectoryView (#1963) * [FIX] Border style * [FIX] Refactoring * [FIX] fix color of border * Undo Co-authored-by: Aroo <azhaubassar@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Clear settings on server change (#1967) * [FIX] Deeplinking without RoomId (#1925) * [FIX] Deeplinking without rid * [FIX] Join channel * [FIX] Deep linking without rid * Update app/lib/methods/canOpenRoom.js Co-authored-by: Diego Mello <diegolmello@gmail.com> * [NEW] Two Factor authentication via email (#1961) * First api call working * [NEW] REST API Post wrapper 2FA * [NEW] Send 2FA on Email * [I18n] Add translations * [NEW] Translations & Cancel totp * [CHORE] Totp -> TwoFactor * [NEW] Two Factor by email * [NEW] Tablet Support * [FIX] Text colors * [NEW] Password 2fa * [FIX] Encrypt password on 2FA * [NEW] MethodCall2FA * [FIX] Password fallback * [FIX] Wrap all post/methodCall with 2fa * [FIX] Wrap missed function * few fixes * [FIX] Use new TOTP on Login * [improvement] 2fa methodCall Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [FIX] Correct message for manual approval user Registration (#1906) * [FIX] Correct message for manual approval from admin shown on Registeration * lint fix - added semicolon * Updated the translations * [FIX] Translations * i18n to match server Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Direct Message between multiple users REST (#1974) * [FIX] Investigate app losing connection issues (#1890) * [WIP] Reopen without timeOut & ping with 5 sec & Fix Unsubscribe * [FIX] Remove duplicated close * [FIX] Use no-dist lib * [FIX] Try minor fix * [FIX] Try reopen connection when app was put on foreground * [FIX] Remove timeout * [FIX] Build * [FIX] Patch * [FIX] Snapshot * [IMPROVEMENT] Decrease time to reopen * [FIX] Some fixes * [FIX] Update sdk version * [FIX] Subscribe Room Once * [CHORE] Update sdk * [FIX] Subscribe Room * [FIX] Try to resend missed subs * [FIX] Users never show status when start app without network * [FIX] Subscribe to room * [FIX] Multiple servers * [CHORE] Update SDK * [FIX] Don't duplicate streams on subscribeAll * [FIX] Server version when start the app offline * [FIX] Server version cached * [CHORE] Remove unnecessary code * [FIX] Offline server version * [FIX] Subscribe before connect * [FIX] Remove unncessary props * [FIX] Update sdk * [FIX] User status & Unsubscribe Typing * [FIX] Typing at incorrect room * [FIX] Multiple Servers * [CHORE] Update SDK * [REVERT] Undo some changes on SDK * [CHORE] Update sdk to prevent incorrect subscribes * [FIX] Prevent no reconnect * [FIX] Remove close on open * [FIX] Clear typing when disconnect/connect to SDK * [CHORE] Update SDK * [CHORE] Update SDK * Update SDK * fix merge develop Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Single message thread inserting thread without rid (#1999) * [FIX] ThreadMessagesView crashing on load (#1997) * [FIX] Saml (#1996) * [FIX] SAML incorrect close * [FIX] Pathname Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Change user own status (#1995) * [FIX] Change user own status * [IMPROVEMENT] Set activeUsers Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Loading all updated rooms after app resume (#1998) * [FIX] Loading all updated rooms after app resume * Fix room date on RoomItem Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Change notifications preferences (#2000) * [FIX] Change notifications preferences * [IMPROVEMENT] Picker View * [I18N] Translations * [FIX] Picker Selection * [FIX] List border * [FIX] Prevent crash * [FIX] Not-Pref tablet * [FIX] Use same style of LanguageView * [IMPROVEMENT] Send listItem title Co-authored-by: Diego Mello <diegolmello@gmail.com> * Bump version to 4.6.1 (#2001) * [FIX] DM header blink (#2011) * [FIX] Split get settings into two requests (#2017) * [FIX] Split get settings into two requests * [FIX] Clear settings only when change server * [IMPROVEMENT] Move the way to clear settings * [REVERT] Revert some changes * [FIX] Server Icon Co-authored-by: Diego Mello <diegolmello@gmail.com> * [REGRESSION] Invite Links (#2007) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Read only channel/broadcast (#1951) * [FIX] Read only channel/broadcast * [FIX] Roles missing * [FIX] Check roles to readOnly * [FIX] Can post * [FIX] Respect post-readonly permission * [FIX] Search a room readOnly Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Cas auth (#2024) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Login TOTP Compatibility to older servers (#2018) * [FIX] Login TOTP Compatibility to older servers * [FIX] Android crashes if use double negation Co-authored-by: Diego Mello <diegolmello@gmail.com> * Bump version to 4.6.4 (#2029) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Lint (#2030) * [FIX] UIKit with only one block (#2022) * [FIX] Message with only one block * [FIX] Update headers Co-authored-by: Diego Mello <diegolmello@gmail.com> * Bump version to 4.7.0 (#2035) * [FIX] Action Tint Color on Black theme (#2081) * [FIX] Prevent crash when thread is not found (#2080) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Prevent double click (#2079) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Show slash commands when disconnected (#2078) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Backhandler onboarding (#2077) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Respect UI_Allow_room_names_with_special_chars setting (#2076) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] RoomsList update sometimes isn't fired (#2071) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Stop inserting last message as message object from rooms stream if room is focused (#2069) * [IMPROVEMENT] No insert last message if the room is focused * fix discussion/threads Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Hide system messages (#2067) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Pending update (#2066) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Prevent crash when room.uids was not inserted yet (#2055) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FEATURE] Save video (#2063) * added-feature-save-video * fix sha256 Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Send totp-code to meteor call (#2050) * fixed-issue * removed-variable-name-errors * reverted-last-commit Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] MessageBox mention shouldn't show group DMs (#2049) * fixed-issue * [FIX] Filter users only if it's not a group chat Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] AttachmentView (Android)(Tablet) (#2047) * [fix]Tablet attachment View and Room Navigation * fix weird navigation and margin bottom Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Allow special chars in Filename (#2020) * fixed-filename-issue * improve Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Recorded audio on Android doesn't play on iOS (#2073) * react-native-video -> expo-av * remove react-native-video * Add audio mode * update mocks * [FIX] Loading bigger than play/pause Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Message Touchable (#2082) * [FIX] Avatar touchable * [IMPROVEMENT] onLongPress on all Message Touchables * [IMPROVEMENT] User & baseUrl on MessageContext * [FIX] Context Access * [FIX] BaseURL * Fix User Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] ReactionsModal (#2085) * [NEW] Delete Server (#1975) * [NEW] Delete server Co-authored-by: Bruno Dantas <oliveiradantas96@gmail.com> Co-authored-by: Calebe Rios <calebersmendes@gmail.com> * [FIX] Revert removed function Co-authored-by: Bruno Dantas <oliveiradantas96@gmail.com> Co-authored-by: Calebe Rios <calebersmendes@gmail.com> * pods * i18n * Revert "pods" This reverts commit 2854a1650538159aeeafe90fdb2118d12b76a82f. Co-authored-by: Bruno Dantas <oliveiradantas96@gmail.com> Co-authored-by: Calebe Rios <calebersmendes@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Change server while connecting/updating (#1981) * [IMPROVEMENT] Change server while connecting * [FIX] Not login/reconnect to previous server * [FIX] Abort all fetch while connecting * [FIX] Abort sdk fetch * [FIX] Patch-package * Add comments Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Keep screen awake while recording/playing some audio (#2089) * [IMPROVEMENT] Keep screen awake while recording/playing some audio * [FIX] Add expo-keep-awake mock * [FIX] UIKit crashing when UIKitModal receive update event (#2088) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Close announcement banner (#2064) * [NEW] Created new field in subscription table Signed-off-by: Ezequiel De Oliveira <ezequiel1de1oliveira@gmail.com> * [NEW] New field added to obeserver in room view Signed-off-by: Ezequiel De Oliveira <ezequiel1de1oliveira@gmail.com> * [NEW] Added icon and new design to banner Signed-off-by: Ezequiel De Oliveira <ezequiel1de1oliveira@gmail.com> * [NEW] Close banner function works Signed-off-by: Ezequiel De Oliveira <ezequiel1de1oliveira@gmail.com> * [IMPROVEMENT] closed banner status now update correctly Signed-off-by: Ezequiel De Oliveira <ezequiel1de1oliveira@gmail.com> * improve banner style Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * Update all dependencies (#2008) * Android RN 62 * First steps iOS * Second step iOS * iOS compiling * "New" build system * Finish iOS * Flipper * Update to RN 0.62.1 * expo libs * Hermes working * Fix lint * Fix android build * Patches * Dev patches * Patch WatermelonDB: https://github.com/Nozbe/WatermelonDB/pull/660 * Fix jitsi * Update several minors * Update dev minors and lint * react-native-keyboard-input * Few updates * device info * react-native-fast-image * Navigation bar color * react-native-picker-select * webview * reactotron-react-native * Watermelondb * RN 0.62.2 * Few updates * Fix selection * update gems * remove lib * finishing * tests * Use node 10 * Re-enable app bundle * iOS build * Update jitsi ios * [NEW] Passcode and biometric unlock (#2059) * Update expo libs * Configure expo-local-authentication * ScreenLockedView * Authenticate server change * Auth on app resume * localAuthentication util * Add servers.lastLocalAuthenticatedSession column * Save last session date on background * Use our own version of app state redux * Fix libs * Remove inactive * ScreenLockConfigView * Apply on saved data * Auto lock option label * Starting passcode * Basic passcode flow working * Change passcode * Check if biometry is enrolled * Use fork * Migration * Patch expo-local-authentication * Use async storage * Styling * Timer * Refactor * Lock orientation portrait when not on tablet * share extension * Deep linking * Share extension * Refactoring passcode * use state * Stash * Refactor * Change passcode * Animate dots on error * Matching passcodes * Shake * Remove lib * Delete button * Fade animation on modal * Refactoring * ItemInfo * I18n * I18n * Remove unnecessary prop * Save biometry column * Raise time to lock to 30 seconds * Vibrate on wrong confirmation passcode * Reset attempts and save last authentication on local passcode confirmation * Remove inline style * Save last auth * Fix header blink * Change function name * Fix android modal * Fix vibration permission * PasscodeEnter calls biometry * Passcode on the state * Biometry button on PasscodeEnter * Show whole passcode * Secure passcode * Save passcode with promise to prevent empty passcodes and immediately lock * Patch expo-local-authentication * I18n * Fix biometry being called every time * Blur screen on app inactive * Revert "Blur screen on app inactive" This reverts commit a4ce812934adcf6cf87eb1a92aec9283e2f26753. * Remove immediately because of how Activities work on Android * Pods * New layout * stash * Layout refactored * Fix icons * Force set passcode from server * Lint * Improve permission message * Forced passcode subtitle * Disable based on admin's choice * Require local authentication on login success * Refactor * Update tests * Update react-native-device-info to fix notch * Lint * Fix modal * Fix icons * Fix min auto lock time * Review * keep enabled on mobile * fix forced by admin when enable unlock with passcode * use DEFAULT_AUTO_LOCK when manual enable screenLock * fix check has passcode * request biometry on first password * reset auto time lock when disabled on server Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> * [FIX] Messages View (#2090) * [FIX] Messages View * [FIX] Opening PDF from Files View * [FIX] Audio * [FIX] SearchMessagesView Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Big names overflow (#2072) * [FIX] Big names overflow * [FIX] Message time Co-authored-by: devyaniChoubey <devyanichoubey16@gmail.com> * [FIX] Some alignments * fix user item overflow * some adjustments Co-authored-by: devyaniChoubey <devyanichoubey16@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Avatar of message as an emoji (#2038) * fixed-issue * removed-hardcoded-emoji * Merge develop * replaced markdown with emoji componenent * made-changes * use avatar onPress Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Diego Mello <diegolmello@gmail.com> * [NEW] Livechat (#2004) * [WIP][NEW] Livechat info/actions * [IMPROVEMENT] RoomActionsView * [NEW] Visitor Navigation * [NEW] Get Department REST * [FIX] Borders * [IMPROVEMENT] Refactor RoomInfo View * [FIX] Error while navigate from mention -> roomInfo * [NEW] Livechat Fields * [NEW] Close Livechat * [WIP] Forward livechat * [NEW] Return inquiry * [WIP] Comment when close livechat * [WIP] Improve roomInfo * [IMPROVEMENT] Forward room * [FIX] Department picker * [FIX] Picker without results * [FIX] Superfluous argument * [FIX] Check permissions on RoomActionsView * [FIX] Livechat permissions * [WIP] Show edit to livechat * [I18N] Add pt-br translations * [WIP] Livechat Info * [IMPROVEMENT] Livechat info * [WIP] Livechat Edit * [WIP] Livechat edit * [WIP] Livechat Edit * [WIP] Livechat edit scroll * [FIX] Edit customFields * [FIX] Clean livechat customField * [FIX] Visitor Navigation * [NEW] Next input logic LivechatEdit * [FIX] Add livechat data to subscription * [FIX] Revert change * [NEW] Livechat user Status * [WIP] Livechat tags * [NEW] Edit livechat tags * [FIX] Prevent some crashes * [FIX] Forward * [FIX] Return Livechat error * [FIX] Prevent livechat info crash * [IMPROVEMENT] Use input style on forward chat * OnboardingSeparator -> OrSeparator * [FIX] Go to next input * [NEW] Added some icons * [NEW] Livechat close * [NEW] Forward Room Action * [FIX] Livechat edit style * [FIX] Change status logic * [CHORE] Remove unnecessary logic * [CHORE] Remove unnecessary code * [CHORE] Remove unecessary case * [FIX] Superfluous argument * [IMPROVEMENT] Submit livechat edit * [CHORE] Remove textInput type * [FIX] Livechat edit * [FIX] Livechat Edit * [FIX] Use same effect * [IMPROVEMENT] Tags input * [FIX] Add empty tag * Fix minor issues * Fix typo * insert livechat room data to our room object * review * add method calls server version Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Delete Subs (#2091) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Android build (#2094) * [FIX] Blink header DM (#2093) * [FIX] Blink header DM * Remove query * [FIX] Push RoomInfoView * remove unnecessary try/catch * [FIX] RoomInfo > Message (Tablet) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Default biometry enabled (#2095) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [IMPROVEMENT] Enable navigating to a room from auth deep linking (#2115) * Wait for login success to navigate * Enable auth and room deep linking at the same time * [FIX] NewMessageView Press Item should open DM (#2116) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Roles throwing error (#2110) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Wait attach activity before changeNavigationBarColor (#2111) * [FIX] Wait attach activity before changeNavigationBarColor * Remove timeout and add try/catch Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] UIKit crash when some app send a list (#2117) * [FIX] StoryBook * [FIX] UIKit crash when some app send a list * [CHORE] Update snapshot * [CHORE] Remove token & id * [FIX] Change bar color while no activity attached (#2130) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Screen Lock options i18n (#2120) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [i18n] Added missing German translation strings (#2105) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Sometimes SDK is null when try to connect (#2131) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [FIX] Autocomplete position on Android (#2106) * [FIX] Autocomplete position on Android * [FIX] Set selection to 0 when needed Co-authored-by: Diego Mello <diegolmello@gmail.com> * Revert "[FIX] Autocomplete position on Android (#2106)" (#2136) This reverts commit e8c38d6f6f69ae396a4aae6e37336617da739a6d. * [FIX] Here and all mentions shouldn't refer to users (#2137) * [FIX] No send data to bugsnag if it's an aborted request (#2133) Co-authored-by: Diego Mello <diegolmello@gmail.com> * [TESTS] Update and separate E2E tests (#2126) * Tests passing until roomslist * create room * roominfo * change server * broadcast * profile * custom status * forgot password * working * room and onboarding * Tests separated * config.yml refactor * Revert "config.yml refactor" This reverts commit 0e984d3029e47612726bf199553f7abdf24843e5. * CI * lint * CI refactor * Onboarding tests * npx detox * Add all tests * Save brew cache * mac-env executor * detox-test command * Update readme * Remove folder * [FIX] Screen Lock Time respect local value (#2141) * [FIX] Screen Lock Time respect local value * [FIX] Enable biometry at the first passcode change Co-authored-by: phriedrich <info@phriedrich.de> Co-authored-by: Guilherme Siqueira <guilhersiqueira@gmail.com> Co-authored-by: Prateek Jain <44807945+Prateek93a@users.noreply.github.com> Co-authored-by: Djorkaeff Alexandre <djorkaeff.unb@gmail.com> Co-authored-by: Prateek Jain <prateek93a@gmail.com> Co-authored-by: devyaniChoubey <52153085+devyaniChoubey@users.noreply.github.com> Co-authored-by: Bernard Seow <ssbing99@gmail.com> Co-authored-by: Hiroki Ishiura <ishiura@ja2.so-net.ne.jp> Co-authored-by: Exordian <jakob.englisch@gmail.com> Co-authored-by: Daanchaam <daanhendriks97@gmail.com> Co-authored-by: Youssef Muhamad <emaildeyoussefmuhamad@gmail.com> Co-authored-by: Iván Álvarez <ialvarezpereira@gmail.com> Co-authored-by: Sarthak Pranesh <41206172+sarthakpranesh@users.noreply.github.com> Co-authored-by: Michele Pellegrini <pellettiero@users.noreply.github.com> Co-authored-by: Tanmoy Bhowmik <tanmoy.openroot@gmail.com> Co-authored-by: Hibikine Kage <14365761+hibikine@users.noreply.github.com> Co-authored-by: Ezequiel de Oliveira <ezequiel1de1oliveira@gmail.com> Co-authored-by: Neil Agarwal <neil@neilagarwal.me> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> Co-authored-by: Govind Dixit <GOVINDDIXIT93@GMAIL.COM> Co-authored-by: Zhaubassarova Aruzhan <49000079+azhaubassar@users.noreply.github.com> Co-authored-by: Aroo <azhaubassar@gmail.com> Co-authored-by: Sarthak Pranesh <sarthak.pranesh2018@vitstudent.ac.in> Co-authored-by: Siddharth Padhi <padhisiddharth31@gmail.com> Co-authored-by: Bruno Dantas <oliveiradantas96@gmail.com> Co-authored-by: Calebe Rios <calebersmendes@gmail.com> Co-authored-by: devyaniChoubey <devyanichoubey16@gmail.com>
2020-05-25 20:54:27 +00:00
// Copyright (c) Facebook, Inc. and its affiliates.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
#pragma once
#include <cassert>
#include <mutex>
#include <utility>
#include "yarpl/flowable/Flowable.h"
#include "yarpl/flowable/Subscriber.h"
#include "yarpl/flowable/Subscription.h"
#include "yarpl/utils/credits.h"
#include <boost/intrusive/list.hpp>
#include <folly/Executor.h>
#include <folly/Synchronized.h>
#include <folly/functional/Invoke.h>
#include <folly/io/async/EventBase.h>
namespace yarpl {
namespace flowable {
/**
* Base (helper) class for operators. Operators are templated on two types: D
* (downstream) and U (upstream). Operators are created by method calls on an
* upstream Flowable, and are Flowables themselves. Multi-stage pipelines can
* be built: a Flowable heading a sequence of Operators.
*/
template <typename U, typename D>
class FlowableOperator : public Flowable<D> {
protected:
/// An Operator's subscription.
///
/// When a pipeline chain is active, each Flowable has a corresponding
/// subscription. Except for the first one, the subscriptions are created
/// against Operators. Each operator subscription has two functions: as a
/// subscriber for the previous stage; as a subscription for the next one, the
/// user-supplied subscriber being the last of the pipeline stages.
class Subscription : public yarpl::flowable::Subscription,
public BaseSubscriber<U> {
protected:
explicit Subscription(std::shared_ptr<Subscriber<D>> subscriber)
: subscriber_(std::move(subscriber)) {
CHECK(yarpl::atomic_load(&subscriber_));
}
// Subscriber will be provided by the init(Subscriber) call
Subscription() {}
virtual void init(std::shared_ptr<Subscriber<D>> subscriber) {
if (yarpl::atomic_load(&subscriber_)) {
subscriber->onSubscribe(yarpl::flowable::Subscription::create());
subscriber->onError(std::runtime_error("already initialized"));
return;
}
subscriber_ = std::move(subscriber);
}
void subscriberOnNext(D value) {
if (auto subscriber = yarpl::atomic_load(&subscriber_)) {
subscriber->onNext(std::move(value));
}
}
/// Terminates both ends of an operator normally.
void terminate() {
std::shared_ptr<Subscriber<D>> null;
auto subscriber = yarpl::atomic_exchange(&subscriber_, null);
BaseSubscriber<U>::cancel();
if (subscriber) {
subscriber->onComplete();
}
}
/// Terminates both ends of an operator with an error.
void terminateErr(folly::exception_wrapper ew) {
std::shared_ptr<Subscriber<D>> null;
auto subscriber = yarpl::atomic_exchange(&subscriber_, null);
BaseSubscriber<U>::cancel();
if (subscriber) {
subscriber->onError(std::move(ew));
}
}
// Subscription.
void request(int64_t n) override {
BaseSubscriber<U>::request(n);
}
void cancel() override {
std::shared_ptr<Subscriber<D>> null;
auto subscriber = yarpl::atomic_exchange(&subscriber_, null);
BaseSubscriber<U>::cancel();
}
// Subscriber.
void onSubscribeImpl() override {
yarpl::atomic_load(&subscriber_)->onSubscribe(this->ref_from_this(this));
}
void onCompleteImpl() override {
std::shared_ptr<Subscriber<D>> null;
if (auto subscriber = yarpl::atomic_exchange(&subscriber_, null)) {
subscriber->onComplete();
}
}
void onErrorImpl(folly::exception_wrapper ew) override {
std::shared_ptr<Subscriber<D>> null;
if (auto subscriber = yarpl::atomic_exchange(&subscriber_, null)) {
subscriber->onError(std::move(ew));
}
}
private:
/// This subscription controls the life-cycle of the subscriber. The
/// subscriber is retained as long as calls on it can be made. (Note: the
/// subscriber in turn maintains a reference on this subscription object
/// until cancellation and/or completion.)
AtomicReference<Subscriber<D>> subscriber_;
};
};
template <typename U, typename D, typename F, typename EF>
class MapOperator : public FlowableOperator<U, D> {
using Super = FlowableOperator<U, D>;
static_assert(std::is_same<std::decay_t<F>, F>::value, "undecayed");
static_assert(folly::is_invocable_r<D, F, U>::value, "not invocable");
static_assert(
folly::is_invocable_r<
folly::exception_wrapper,
EF,
folly::exception_wrapper&&>::value,
"exception handler not invocable");
public:
template <typename Func, typename ErrorFunc>
MapOperator(
std::shared_ptr<Flowable<U>> upstream,
Func&& function,
ErrorFunc&& errFunction)
: upstream_(std::move(upstream)),
function_(std::forward<Func>(function)),
errFunction_(std::move(errFunction)) {}
void subscribe(std::shared_ptr<Subscriber<D>> subscriber) override {
upstream_->subscribe(std::make_shared<Subscription>(
this->ref_from_this(this), std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(
std::shared_ptr<MapOperator> flowable,
std::shared_ptr<Subscriber<D>> subscriber)
: SuperSubscription(std::move(subscriber)),
flowable_(std::move(flowable)) {}
void onNextImpl(U value) override {
try {
if (auto flowable = yarpl::atomic_load(&flowable_)) {
this->subscriberOnNext(flowable->function_(std::move(value)));
}
} catch (const std::exception& exn) {
folly::exception_wrapper ew{std::current_exception(), exn};
this->terminateErr(std::move(ew));
}
}
void onErrorImpl(folly::exception_wrapper ew) override {
try {
if (auto flowable = yarpl::atomic_load(&flowable_)) {
SuperSubscription::onErrorImpl(flowable->errFunction_(std::move(ew)));
}
} catch (const std::exception& exn) {
this->terminateErr(
folly::exception_wrapper{std::current_exception(), exn});
}
}
void onTerminateImpl() override {
yarpl::atomic_exchange(&flowable_, nullptr);
SuperSubscription::onTerminateImpl();
}
private:
AtomicReference<MapOperator> flowable_;
};
std::shared_ptr<Flowable<U>> upstream_;
F function_;
EF errFunction_;
};
template <typename U, typename F>
class FilterOperator : public FlowableOperator<U, U> {
// for use in subclasses
using Super = FlowableOperator<U, U>;
static_assert(std::is_same<std::decay_t<F>, F>::value, "undecayed");
static_assert(folly::is_invocable_r<bool, F, U>::value, "not invocable");
public:
template <typename Func>
FilterOperator(std::shared_ptr<Flowable<U>> upstream, Func&& function)
: upstream_(std::move(upstream)),
function_(std::forward<Func>(function)) {}
void subscribe(std::shared_ptr<Subscriber<U>> subscriber) override {
upstream_->subscribe(std::make_shared<Subscription>(
this->ref_from_this(this), std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(
std::shared_ptr<FilterOperator> flowable,
std::shared_ptr<Subscriber<U>> subscriber)
: SuperSubscription(std::move(subscriber)),
flowable_(std::move(flowable)) {}
void onNextImpl(U value) override {
if (auto flowable = yarpl::atomic_load(&flowable_)) {
if (flowable->function_(value)) {
SuperSubscription::subscriberOnNext(std::move(value));
} else {
SuperSubscription::request(1);
}
}
}
void onTerminateImpl() override {
yarpl::atomic_exchange(&flowable_, nullptr);
SuperSubscription::onTerminateImpl();
}
private:
AtomicReference<FilterOperator> flowable_;
};
std::shared_ptr<Flowable<U>> upstream_;
F function_;
};
template <typename U, typename D, typename F>
class ReduceOperator : public FlowableOperator<U, D> {
using Super = FlowableOperator<U, D>;
static_assert(std::is_same<std::decay_t<F>, F>::value, "undecayed");
static_assert(std::is_assignable<D&, U>::value, "not assignable");
static_assert(folly::is_invocable_r<D, F, D, U>::value, "not invocable");
public:
template <typename Func>
ReduceOperator(std::shared_ptr<Flowable<U>> upstream, Func&& function)
: upstream_(std::move(upstream)),
function_(std::forward<Func>(function)) {}
void subscribe(std::shared_ptr<Subscriber<D>> subscriber) override {
upstream_->subscribe(std::make_shared<Subscription>(
this->ref_from_this(this), std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(
std::shared_ptr<ReduceOperator> flowable,
std::shared_ptr<Subscriber<D>> subscriber)
: SuperSubscription(std::move(subscriber)),
flowable_(std::move(flowable)),
accInitialized_(false) {}
void request(int64_t) override {
// Request all of the items.
SuperSubscription::request(credits::kNoFlowControl);
}
void onNextImpl(U value) override {
if (accInitialized_) {
if (auto flowable = yarpl::atomic_load(&flowable_)) {
acc_ = flowable->function_(std::move(acc_), std::move(value));
}
} else {
acc_ = std::move(value);
accInitialized_ = true;
}
}
void onCompleteImpl() override {
if (accInitialized_) {
SuperSubscription::subscriberOnNext(std::move(acc_));
}
SuperSubscription::onCompleteImpl();
}
void onTerminateImpl() override {
yarpl::atomic_exchange(&flowable_, nullptr);
SuperSubscription::onTerminateImpl();
}
private:
AtomicReference<ReduceOperator> flowable_;
bool accInitialized_;
D acc_;
};
std::shared_ptr<Flowable<U>> upstream_;
F function_;
};
template <typename T>
class TakeOperator : public FlowableOperator<T, T> {
using Super = FlowableOperator<T, T>;
public:
TakeOperator(std::shared_ptr<Flowable<T>> upstream, int64_t limit)
: upstream_(std::move(upstream)), limit_(limit) {}
void subscribe(std::shared_ptr<Subscriber<T>> subscriber) override {
upstream_->subscribe(
std::make_shared<Subscription>(limit_, std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(int64_t limit, std::shared_ptr<Subscriber<T>> subscriber)
: SuperSubscription(std::move(subscriber)), limit_(limit) {}
void onSubscribeImpl() override {
SuperSubscription::onSubscribeImpl();
if (limit_ <= 0) {
SuperSubscription::terminate();
}
}
void onNextImpl(T value) override {
if (limit_-- > 0) {
if (pending_ > 0) {
--pending_;
}
SuperSubscription::subscriberOnNext(std::move(value));
if (limit_ == 0) {
SuperSubscription::terminate();
}
}
}
void request(int64_t delta) override {
delta = std::min(delta, limit_ - pending_);
if (delta > 0) {
pending_ += delta;
SuperSubscription::request(delta);
}
}
private:
int64_t pending_{0};
int64_t limit_;
};
std::shared_ptr<Flowable<T>> upstream_;
const int64_t limit_;
};
template <typename T>
class SkipOperator : public FlowableOperator<T, T> {
using Super = FlowableOperator<T, T>;
public:
SkipOperator(std::shared_ptr<Flowable<T>> upstream, int64_t offset)
: upstream_(std::move(upstream)), offset_(offset) {}
void subscribe(std::shared_ptr<Subscriber<T>> subscriber) override {
upstream_->subscribe(
std::make_shared<Subscription>(offset_, std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(int64_t offset, std::shared_ptr<Subscriber<T>> subscriber)
: SuperSubscription(std::move(subscriber)), offset_(offset) {}
void onNextImpl(T value) override {
if (offset_ > 0) {
--offset_;
} else {
SuperSubscription::subscriberOnNext(std::move(value));
}
}
void request(int64_t delta) override {
if (firstRequest_) {
firstRequest_ = false;
delta = credits::add(delta, offset_);
}
SuperSubscription::request(delta);
}
private:
int64_t offset_;
bool firstRequest_{true};
};
std::shared_ptr<Flowable<T>> upstream_;
const int64_t offset_;
};
template <typename T>
class IgnoreElementsOperator : public FlowableOperator<T, T> {
using Super = FlowableOperator<T, T>;
public:
explicit IgnoreElementsOperator(std::shared_ptr<Flowable<T>> upstream)
: upstream_(std::move(upstream)) {}
void subscribe(std::shared_ptr<Subscriber<T>> subscriber) override {
upstream_->subscribe(std::make_shared<Subscription>(std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(std::shared_ptr<Subscriber<T>> subscriber)
: SuperSubscription(std::move(subscriber)) {}
void onNextImpl(T) override {}
};
std::shared_ptr<Flowable<T>> upstream_;
};
template <typename T>
class SubscribeOnOperator : public FlowableOperator<T, T> {
using Super = FlowableOperator<T, T>;
public:
SubscribeOnOperator(
std::shared_ptr<Flowable<T>> upstream,
folly::Executor& executor)
: upstream_(std::move(upstream)), executor_(executor) {}
void subscribe(std::shared_ptr<Subscriber<T>> subscriber) override {
executor_.add([this, self = this->ref_from_this(this), subscriber] {
upstream_->subscribe(
std::make_shared<Subscription>(executor_, std::move(subscriber)));
});
}
private:
using SuperSubscription = typename Super::Subscription;
class Subscription : public SuperSubscription {
public:
Subscription(
folly::Executor& executor,
std::shared_ptr<Subscriber<T>> subscriber)
: SuperSubscription(std::move(subscriber)), executor_(executor) {}
void request(int64_t delta) override {
executor_.add([delta, this, self = this->ref_from_this(this)] {
this->callSuperRequest(delta);
});
}
void cancel() override {
executor_.add([this, self = this->ref_from_this(this)] {
this->callSuperCancel();
});
}
void onNextImpl(T value) override {
SuperSubscription::subscriberOnNext(std::move(value));
}
private:
// Trampoline to call superclass method; gcc bug 58972.
void callSuperRequest(int64_t delta) {
SuperSubscription::request(delta);
}
// Trampoline to call superclass method; gcc bug 58972.
void callSuperCancel() {
SuperSubscription::cancel();
}
folly::Executor& executor_;
};
std::shared_ptr<Flowable<T>> upstream_;
folly::Executor& executor_;
};
template <typename T, typename OnSubscribe>
class FromPublisherOperator : public Flowable<T> {
static_assert(
std::is_same<std::decay_t<OnSubscribe>, OnSubscribe>::value,
"undecayed");
public:
template <typename F>
explicit FromPublisherOperator(F&& function)
: function_(std::forward<F>(function)) {}
void subscribe(std::shared_ptr<Subscriber<T>> subscriber) override {
function_(std::move(subscriber));
}
private:
OnSubscribe function_;
};
template <typename T, typename R>
class FlatMapOperator : public FlowableOperator<T, R> {
using Super = FlowableOperator<T, R>;
public:
FlatMapOperator(
std::shared_ptr<Flowable<T>> upstream,
folly::Function<std::shared_ptr<Flowable<R>>(T)> func)
: upstream_(std::move(upstream)), function_(std::move(func)) {}
void subscribe(std::shared_ptr<Subscriber<R>> subscriber) override {
upstream_->subscribe(std::make_shared<FMSubscription>(
this->ref_from_this(this), std::move(subscriber)));
}
private:
using SuperSubscription = typename Super::Subscription;
class FMSubscription : public SuperSubscription {
struct MappedStreamSubscriber;
public:
FMSubscription(
std::shared_ptr<FlatMapOperator> flowable,
std::shared_ptr<Subscriber<R>> subscriber)
: SuperSubscription(std::move(subscriber)),
flowable_(std::move(flowable)) {}
void onSubscribeImpl() final {
liveSubscribers_++;
SuperSubscription::onSubscribeImpl();
}
void onNextImpl(T value) final {
std::shared_ptr<Flowable<R>> mappedStream;
try {
mappedStream = flowable_->function_(std::move(value));
} catch (const std::exception& exn) {
folly::exception_wrapper ew{std::current_exception(), exn};
{
std::lock_guard<std::mutex> g(onErrorExGuard_);
onErrorEx_ = ew;
}
// next iteration of drainLoop will cancel this subscriber as well
drainLoop();
return;
}
std::shared_ptr<MappedStreamSubscriber> mappedSubscriber =
std::make_shared<MappedStreamSubscriber>(this->ref_from_this(this));
mappedSubscriber->fmReference_ = mappedSubscriber;
{
// put into pendingValue queue because once the mappedSubscriber
// is subscribed to, it will request elements. We don't want the
// drainLoop to execute while it's on withoutValue, and request
// a second element before the first arrives.
auto l = lists.wlock();
CHECK(!mappedSubscriber->is_linked());
l->pendingValue.push_back(*mappedSubscriber.get());
}
liveSubscribers_++;
mappedStream->subscribe(mappedSubscriber);
drainLoop();
}
void drainImpl() {
// phase 1: clear out terminated subscribers
{
auto clearList = [](auto& list, SubscriberList& t) {
while (!list.empty()) {
auto& elem = list.front();
auto r = elem.sync.wlock();
r->freeze = true;
elem.unlink();
t.push_back(elem);
}
};
SubscriberList clearTrash;
if (clearAllSubscribers_.load()) {
auto l = lists.wlock();
clearList(l->withValue, clearTrash);
clearList(l->withoutValue, clearTrash);
clearList(l->pendingValue, clearTrash);
}
// clear elements while no locks are held
while (!clearTrash.empty()) {
auto& elem = clearTrash.front();
elem.unlink();
elem.cancel();
elem.fmReference_ = nullptr;
}
}
// phase 2: check if the subscriber should terminate due to error
// or all subscribers completing
if (!calledDownstreamTerminate_) {
folly::exception_wrapper ex;
{
std::lock_guard<std::mutex> exg(onErrorExGuard_);
ex = std::move(onErrorEx_);
}
if (ex) {
calledDownstreamTerminate_ = true;
cancel();
this->terminateErr(std::move(ex));
} else if (liveSubscribers_ == 0) {
calledDownstreamTerminate_ = true;
this->terminate();
}
}
// phase 3: if the downstream has requested elements, pop values out of
// subscribers which have received a value and call downstream->onNext
while (requested_ != 0) {
R val;
{
auto l = lists.wlock();
if (l->withValue.empty()) {
break;
}
requested_--;
auto& elem = l->withValue.front();
elem.unlink();
{
auto r = elem.sync.wlock();
CHECK(r->hasValue);
r->hasValue = false;
val = std::move(r->value);
l->withoutValue.push_back(elem);
}
}
SuperSubscription::subscriberOnNext(std::move(val));
}
// phase 4: ask any upstream flowables which don't have pending
// requests for their next element kick off any more requests.
// Put subscribers which have terminated into the trash.
{
SubscriberList terminatedTrash;
while (true) {
MappedStreamSubscriber* elem;
{
auto l = lists.wlock();
if (l->withoutValue.empty()) {
break;
}
elem = &l->withoutValue.front();
auto r = elem->sync.wlock();
CHECK(!r->hasValue) << "failed for elem=" << elem; // sanity
elem->unlink();
// Subscribers might call onNext and then terminate; delay
// removing its liveSubscriber reference until we've delivered
// its element to the downstream subscriber and dropped its
// synchronized reference to `r`, as dropping the
// flatMapSubscription_ reference may invoke its destructor
if (r->isTerminated) {
r->freeze = true;
terminatedTrash.push_back(*elem);
continue; // skips the next elem->request(1)
}
// else, the stream hasn't terminated, request another
// element
l->pendingValue.push_back(*elem);
}
elem->request(1);
}
// phase 5: destroy any mapped subscribers which have terminated,
// enqueue another drain loop run if we do end up discarding any
// subscribers, as our live subscriber count may have gone to zero
if (!terminatedTrash.empty()) {
drainLoopMutex_++;
}
while (!terminatedTrash.empty()) {
auto& elem = terminatedTrash.front();
CHECK(elem.sync.wlock()->isTerminated);
elem.unlink();
elem.fmReference_ = nullptr;
liveSubscribers_--;
}
}
}
// called from MappedStreamSubscriber, receives the R and the
// subscriber which generated the R
void drainLoop() {
auto self = this->ref_from_this(this);
if (drainLoopMutex_++ == 0) {
do {
drainImpl();
} while (drainLoopMutex_-- != 1);
}
}
void onMappedSubscriberNext(MappedStreamSubscriber* elem, R value) {
{
// `elem` may not be in a list, as it may have been canceled. Push it
// on the withValue list and let drainLoop clear it if that's the case.
auto l = lists.wlock();
auto r = elem->sync.wlock();
if (r->freeze) {
return;
}
CHECK(!r->hasValue) << "failed for elem=" << elem;
r->hasValue = true;
r->value = std::move(value);
elem->unlink();
l->withValue.push_back(*elem);
}
drainLoop();
}
void onMappedSubscriberTerminate(MappedStreamSubscriber* elem) {
{
auto r = elem->sync.wlock();
r->isTerminated = true;
if (r->onErrorEx) {
std::lock_guard<std::mutex> exg(onErrorExGuard_);
onErrorEx_ = std::move(r->onErrorEx);
}
if (r->freeze) {
return;
}
}
{
auto l = lists.wlock();
auto r = elem->sync.wlock();
if (r->freeze) {
return;
}
CHECK(elem->is_linked());
elem->unlink();
if (r->hasValue) {
l->withValue.push_back(*elem);
} else {
liveSubscribers_--;
elem->fmReference_ = nullptr;
}
}
drainLoop();
}
// onComplete/onError fall through to onTerminateImpl, which
// will call drainLoop and update the liveSubscribers_ count
void onCompleteImpl() final {}
void onErrorImpl(folly::exception_wrapper ex) final {
std::lock_guard<std::mutex> g(onErrorExGuard_);
onErrorEx_ = std::move(ex);
clearAllSubscribers_.store(true);
}
void onTerminateImpl() final {
liveSubscribers_--;
drainLoop();
flowable_.reset();
}
void request(int64_t n) override {
if ((n + requested_) < requested_) {
requested_ = std::numeric_limits<int64_t>::max();
} else {
requested_ += n;
}
if (n > 0) {
// TODO: make max parallelism configurable a-la RxJava 2.x's
// FlowableFlatMapOperator
SuperSubscription::request(std::numeric_limits<int64_t>::max());
}
drainLoop();
}
void cancel() override {
clearAllSubscribers_.store(true);
drainLoop();
}
private:
// buffers at most a single element of type R
struct MappedStreamSubscriber
: public BaseSubscriber<R>,
public boost::intrusive::list_base_hook<
boost::intrusive::link_mode<boost::intrusive::auto_unlink>> {
MappedStreamSubscriber(std::shared_ptr<FMSubscription> subscription)
: flatMapSubscription_(std::move(subscription)) {}
void onSubscribeImpl() final {
auto fmsb = yarpl::atomic_load(&flatMapSubscription_);
if (!fmsb || fmsb->clearAllSubscribers_) {
BaseSubscriber<R>::cancel();
return;
}
#ifndef NDEBUG
if (auto fms = yarpl::atomic_load(&flatMapSubscription_)) {
auto l = fms->lists.wlock();
auto r = sync.wlock();
if (!is_in_list(*this, l->pendingValue, l)) {
LOG(INFO) << "failed: this=" << this;
LOG(INFO) << "in list: ";
debug_is_in_list(*this, l);
DCHECK(r->freeze);
} else {
}
DCHECK(!r->hasValue);
}
#endif
BaseSubscriber<R>::request(1);
}
void onNextImpl(R value) final {
if (auto fms = yarpl::atomic_load(&flatMapSubscription_)) {
fms->onMappedSubscriberNext(this, std::move(value));
}
}
// noop
void onCompleteImpl() final {}
void onErrorImpl(folly::exception_wrapper ex) final {
auto r = sync.wlock();
r->onErrorEx = std::move(ex);
}
void onTerminateImpl() override {
std::shared_ptr<FMSubscription> null;
if (auto fms = yarpl::atomic_exchange(&flatMapSubscription_, null)) {
fms->onMappedSubscriberTerminate(this);
}
}
struct SyncData {
R value;
bool hasValue{false};
bool isTerminated{false};
bool freeze{false};
folly::exception_wrapper onErrorEx{nullptr};
};
folly::Synchronized<SyncData> sync;
// FMSubscription's 'reference' to this object. FMSubscription
// clears this reference when it drops the MappedStreamSubscriber
// from one of its atomic lists
std::shared_ptr<MappedStreamSubscriber> fmReference_{nullptr};
// this is both a Subscriber and a Subscription<T>
AtomicReference<FMSubscription> flatMapSubscription_{nullptr};
};
// used to make sure only one thread at a time is calling subscriberOnNext
std::atomic<int64_t> drainLoopMutex_{0};
using SubscriberList = boost::intrusive::list<
MappedStreamSubscriber,
boost::intrusive::constant_time_size<false>>;
struct Lists {
// subscribers with a ready R
SubscriberList withValue{};
// subscribers that have requested 1 R, waiting for it to arrive via
// onNext
SubscriberList pendingValue{};
// idle subscribers
SubscriberList withoutValue{};
};
folly::Synchronized<Lists> lists;
template <typename L>
static bool is_in_list(
MappedStreamSubscriber const& elem,
SubscriberList const& list,
L const& lists) {
return in_list_impl(elem, list, lists, true);
}
template <typename L>
static bool not_in_list(
MappedStreamSubscriber const& elem,
SubscriberList const& list,
L const& lists) {
return in_list_impl(elem, list, lists, false);
}
template <typename L>
static bool in_list_impl(
MappedStreamSubscriber const& elem,
SubscriberList const& list,
L const& lists,
bool should) {
if (is_in_list(elem, list) != should) {
#ifndef NDEBUG
debug_is_in_list(elem, lists);
#else
(void)lists;
#endif
return false;
}
return true;
}
template <typename L>
static void debug_is_in_list(
MappedStreamSubscriber const& elem,
L const& lists) {
LOG(INFO) << "in without: " << is_in_list(elem, lists->withoutValue);
LOG(INFO) << "in pending: " << is_in_list(elem, lists->pendingValue);
LOG(INFO) << "in withval: " << is_in_list(elem, lists->withValue);
}
static bool is_in_list(
MappedStreamSubscriber const& elem,
SubscriberList const& list) {
bool found = false;
for (auto& e : list) {
if (&e == &elem) {
found = true;
break;
}
}
return found;
}
std::shared_ptr<FlatMapOperator> flowable_;
// got a terminating signal from the upstream flowable
// always modified in the protected drainImpl()
bool calledDownstreamTerminate_{false};
std::mutex onErrorExGuard_;
folly::exception_wrapper onErrorEx_{nullptr};
// clear all lists of
std::atomic<bool> clearAllSubscribers_{false};
std::atomic<int64_t> requested_{0};
// number of subscribers (FMSubscription + MappedStreamSubscriber) which
// have not received a terminating signal yet
std::atomic<int64_t> liveSubscribers_{0};
};
std::shared_ptr<Flowable<T>> upstream_;
folly::Function<std::shared_ptr<Flowable<R>>(T)> function_;
};
} // namespace flowable
} // namespace yarpl
#include "yarpl/flowable/FlowableConcatOperators.h"
#include "yarpl/flowable/FlowableDoOperator.h"
#include "yarpl/flowable/FlowableObserveOnOperator.h"
#include "yarpl/flowable/FlowableTimeoutOperator.h"