From 876ee25e583727fc8afaf9d3ce7d7a07e845c46a Mon Sep 17 00:00:00 2001 From: Sergei Sysoev Date: Thu, 30 Jan 2025 02:55:58 +0100 Subject: [PATCH] [fleet] Fix multiplatform compilation errors GitOrigin-RevId: 3a62de10db0f4857a851d2e1c4df91e9cdac8264 --- .../andel/intervals/impl/MergingIterator.kt | 4 +- fleet/andel/src/andel/text/CodepointClass.kt | 138 ++++++++++++++- fleet/kernel/src/fleet/kernel/ReteExt.kt | 4 +- .../src/fleet/kernel/rebase/RebaseLoop.kt | 2 +- .../src/fleet/kernel/rete/WithMatchResult.kt | 8 +- .../src/fleet/kernel/rete/impl/Tracing.kt | 2 +- fleet/rpc/src/fleet/rpc/client/RpcClient.kt | 167 +++++++++--------- fleet/rpc/src/fleet/rpc/core/RpcStream.kt | 5 +- 8 files changed, 232 insertions(+), 98 deletions(-) diff --git a/fleet/andel/src/andel/intervals/impl/MergingIterator.kt b/fleet/andel/src/andel/intervals/impl/MergingIterator.kt index a3d2078ad9b7..4a62acd78e6e 100644 --- a/fleet/andel/src/andel/intervals/impl/MergingIterator.kt +++ b/fleet/andel/src/andel/intervals/impl/MergingIterator.kt @@ -37,7 +37,7 @@ internal class MergingIterator(private var first: IntervalsIterator, val firstNext = first.next() val secondNext = second!!.next() if (firstNext && secondNext) { - if (comparator.compare(first, second) > 0) { + if (comparator.compare(first, second!!) > 0) { val tmp = second second = first first = tmp!! @@ -60,7 +60,7 @@ internal class MergingIterator(private var first: IntervalsIterator, else { if (first.next()) { if (second != null) { - if (comparator.compare(first, second) > 0) { + if (comparator.compare(first, second!!) > 0) { val tmp: IntervalsIterator = second as IntervalsIterator second = first first = tmp diff --git a/fleet/andel/src/andel/text/CodepointClass.kt b/fleet/andel/src/andel/text/CodepointClass.kt index 51cae35542be..69928af8aafe 100644 --- a/fleet/andel/src/andel/text/CodepointClass.kt +++ b/fleet/andel/src/andel/text/CodepointClass.kt @@ -17,8 +17,8 @@ fun codepointClass(codepoint: Int): CodepointClass = codepoint == '\r'.code -> CodepointClass.NEWLINE codepoint == '_'.code -> CodepointClass.UNDERSCORE codepoint <= maxSeparatorCode && separatorCodes[codepoint] -> CodepointClass.SEPARATOR - Character.isWhitespace(codepoint) -> CodepointClass.SPACE - Character.isUpperCase(codepoint) -> CodepointClass.UPPERCASE + isWhitespace(codepoint) -> CodepointClass.SPACE + isUpperCase(codepoint) -> CodepointClass.UPPERCASE else -> CodepointClass.LOWERCASE // treat all lowercase and unicode symbols as lowercase } @@ -34,6 +34,56 @@ internal fun charGeomLength(char: Char): Float { } } +// Code points are derived from: +// https://unicode.org/Public/UNIDATA/PropList.txt +fun isWhitespace(codepoint: Int): Boolean { + return codepoint in 0x0009..0x000d // .. + || codepoint == 0x0020 // SPACE + || codepoint == 0x0085 // + || codepoint == 0x00a0 // NO-BREAK SPACE + || codepoint > 0x1000 && + (codepoint == 0x1680 // OGHAM SPACE MARK + || codepoint in 0x2000..0x200a // EN QUAD..HAIR SPACE + || codepoint == 0x2028 // LINE SEPARATOR + || codepoint == 0x2029 // PARAGRAPH SEPARATOR + || codepoint == 0x202f // NARROW NO-BREAK SPACE + || codepoint == 0x205f // MEDIUM MATHEMATICAL SPACE + || codepoint == 0x3000) // IDEOGRAPHIC SPACE +} + +fun isUpperCase(codepoint: Int): Boolean { + // fast path + if (codepoint in 'A'.code..'Z'.code) { + return true + } + if (codepoint < '\u0080'.code) { + return false + } + // proper check + return isLuCategory(codepoint) || isOtherUppercase(codepoint) +} + +private fun isLuCategory(codepoint: Int): Boolean { + val insertionPoint = upperCaseRangeStarts.binarySearch(codepoint) + return when { + insertionPoint >= 0 -> true + insertionPoint == -1 -> false + else -> { + codepoint <= upperCaseRangeEnds[-(insertionPoint + 2)] + } + } +} + +// Code points are derived from: +// https://unicode.org/Public/UNIDATA/PropList.txt +private fun isOtherUppercase(codepoint: Int): Boolean { + return codepoint in 0x2160..0x216f // ROMAN NUMERAL ONE..ROMAN NUMERAL ONE THOUSAND + || codepoint in 0x24b6..0x24cf // CIRCLED LATIN CAPITAL LETTER A..CIRCLED LATIN CAPITAL LETTER Z + || codepoint in 0x1f130..0x1f149 // SQUARED LATIN CAPITAL LETTER A..SQUARED LATIN CAPITAL LETTER Z + || codepoint in 0x1f150..0x1f169 // NEGATIVE CIRCLED LATIN CAPITAL LETTER A..NEGATIVE CIRCLED LATIN CAPITAL LETTER Z + || codepoint in 0x1f170..0x1f189 // NEGATIVE SQUARED LATIN CAPITAL LETTER A..NEGATIVE SQUARED LATIN CAPITAL LETTER Z +} + // Code points are derived from: // https://unicode.org/Public/UNIDATA/EastAsianWidth.txt fun isFullWidth(codePoint: Int): Boolean { @@ -238,6 +288,90 @@ private val AMBIGUOUS = arrayOf(charArrayOf(0x00A1.toChar(), 0x00A1.toChar()), c charArrayOf(0x273D.toChar(), 0x273D.toChar()), charArrayOf(0x2776.toChar(), 0x277F.toChar()), charArrayOf(0xE000.toChar(), 0xF8FF.toChar()), charArrayOf(0xFFFD.toChar(), 0xFFFD.toChar())) +// extracted from https://unicode.org/Public/UNIDATA/UnicodeData.txt +private val upperCaseRangeStarts = intArrayOf( + 0x0041, 0x00c0, 0x00d8, 0x0100, 0x0102, 0x0104, 0x0106, 0x0108, 0x010a, 0x010c, 0x010e, 0x0110, 0x0112, 0x0114, 0x0116, 0x0118, 0x011a, + 0x011c, 0x011e, 0x0120, 0x0122, 0x0124, 0x0126, 0x0128, 0x012a, 0x012c, 0x012e, 0x0130, 0x0132, 0x0134, 0x0136, 0x0139, 0x013b, 0x013d, + 0x013f, 0x0141, 0x0143, 0x0145, 0x0147, 0x014a, 0x014c, 0x014e, 0x0150, 0x0152, 0x0154, 0x0156, 0x0158, 0x015a, 0x015c, 0x015e, 0x0160, + 0x0162, 0x0164, 0x0166, 0x0168, 0x016a, 0x016c, 0x016e, 0x0170, 0x0172, 0x0174, 0x0176, 0x0178, 0x017b, 0x017d, 0x0181, 0x0184, 0x0186, + 0x0189, 0x018e, 0x0193, 0x0196, 0x019c, 0x019f, 0x01a2, 0x01a4, 0x01a6, 0x01a9, 0x01ac, 0x01ae, 0x01b1, 0x01b5, 0x01b7, 0x01bc, 0x01c4, + 0x01c7, 0x01ca, 0x01cd, 0x01cf, 0x01d1, 0x01d3, 0x01d5, 0x01d7, 0x01d9, 0x01db, 0x01de, 0x01e0, 0x01e2, 0x01e4, 0x01e6, 0x01e8, 0x01ea, + 0x01ec, 0x01ee, 0x01f1, 0x01f4, 0x01f6, 0x01fa, 0x01fc, 0x01fe, 0x0200, 0x0202, 0x0204, 0x0206, 0x0208, 0x020a, 0x020c, 0x020e, 0x0210, + 0x0212, 0x0214, 0x0216, 0x0218, 0x021a, 0x021c, 0x021e, 0x0220, 0x0222, 0x0224, 0x0226, 0x0228, 0x022a, 0x022c, 0x022e, 0x0230, 0x0232, + 0x023a, 0x023d, 0x0241, 0x0243, 0x0248, 0x024a, 0x024c, 0x024e, 0x0370, 0x0372, 0x0376, 0x037f, 0x0386, 0x0388, 0x038c, 0x038e, 0x0391, + 0x03a3, 0x03cf, 0x03d2, 0x03d8, 0x03da, 0x03dc, 0x03de, 0x03e0, 0x03e2, 0x03e4, 0x03e6, 0x03e8, 0x03ea, 0x03ec, 0x03ee, 0x03f4, 0x03f7, + 0x03f9, 0x03fd, 0x0460, 0x0462, 0x0464, 0x0466, 0x0468, 0x046a, 0x046c, 0x046e, 0x0470, 0x0472, 0x0474, 0x0476, 0x0478, 0x047a, 0x047c, + 0x047e, 0x0480, 0x048a, 0x048c, 0x048e, 0x0490, 0x0492, 0x0494, 0x0496, 0x0498, 0x049a, 0x049c, 0x049e, 0x04a0, 0x04a2, 0x04a4, 0x04a6, + 0x04a8, 0x04aa, 0x04ac, 0x04ae, 0x04b0, 0x04b2, 0x04b4, 0x04b6, 0x04b8, 0x04ba, 0x04bc, 0x04be, 0x04c0, 0x04c3, 0x04c5, 0x04c7, 0x04c9, + 0x04cb, 0x04cd, 0x04d0, 0x04d2, 0x04d4, 0x04d6, 0x04d8, 0x04da, 0x04dc, 0x04de, 0x04e0, 0x04e2, 0x04e4, 0x04e6, 0x04e8, 0x04ea, 0x04ec, + 0x04ee, 0x04f0, 0x04f2, 0x04f4, 0x04f6, 0x04f8, 0x04fa, 0x04fc, 0x04fe, 0x0500, 0x0502, 0x0504, 0x0506, 0x0508, 0x050a, 0x050c, 0x050e, + 0x0510, 0x0512, 0x0514, 0x0516, 0x0518, 0x051a, 0x051c, 0x051e, 0x0520, 0x0522, 0x0524, 0x0526, 0x0528, 0x052a, 0x052c, 0x052e, 0x0531, + 0x10a0, 0x10c7, 0x10cd, 0x13a0, 0x1c89, 0x1c90, 0x1cbd, 0x1e00, 0x1e02, 0x1e04, 0x1e06, 0x1e08, 0x1e0a, 0x1e0c, 0x1e0e, 0x1e10, 0x1e12, + 0x1e14, 0x1e16, 0x1e18, 0x1e1a, 0x1e1c, 0x1e1e, 0x1e20, 0x1e22, 0x1e24, 0x1e26, 0x1e28, 0x1e2a, 0x1e2c, 0x1e2e, 0x1e30, 0x1e32, 0x1e34, + 0x1e36, 0x1e38, 0x1e3a, 0x1e3c, 0x1e3e, 0x1e40, 0x1e42, 0x1e44, 0x1e46, 0x1e48, 0x1e4a, 0x1e4c, 0x1e4e, 0x1e50, 0x1e52, 0x1e54, 0x1e56, + 0x1e58, 0x1e5a, 0x1e5c, 0x1e5e, 0x1e60, 0x1e62, 0x1e64, 0x1e66, 0x1e68, 0x1e6a, 0x1e6c, 0x1e6e, 0x1e70, 0x1e72, 0x1e74, 0x1e76, 0x1e78, + 0x1e7a, 0x1e7c, 0x1e7e, 0x1e80, 0x1e82, 0x1e84, 0x1e86, 0x1e88, 0x1e8a, 0x1e8c, 0x1e8e, 0x1e90, 0x1e92, 0x1e94, 0x1e9e, 0x1ea0, 0x1ea2, + 0x1ea4, 0x1ea6, 0x1ea8, 0x1eaa, 0x1eac, 0x1eae, 0x1eb0, 0x1eb2, 0x1eb4, 0x1eb6, 0x1eb8, 0x1eba, 0x1ebc, 0x1ebe, 0x1ec0, 0x1ec2, 0x1ec4, + 0x1ec6, 0x1ec8, 0x1eca, 0x1ecc, 0x1ece, 0x1ed0, 0x1ed2, 0x1ed4, 0x1ed6, 0x1ed8, 0x1eda, 0x1edc, 0x1ede, 0x1ee0, 0x1ee2, 0x1ee4, 0x1ee6, + 0x1ee8, 0x1eea, 0x1eec, 0x1eee, 0x1ef0, 0x1ef2, 0x1ef4, 0x1ef6, 0x1ef8, 0x1efa, 0x1efc, 0x1efe, 0x1f08, 0x1f18, 0x1f28, 0x1f38, 0x1f48, + 0x1f59, 0x1f5b, 0x1f5d, 0x1f5f, 0x1f68, 0x1fb8, 0x1fc8, 0x1fd8, 0x1fe8, 0x1ff8, 0x2102, 0x2107, 0x210b, 0x2110, 0x2115, 0x2119, 0x2124, + 0x2126, 0x2128, 0x212a, 0x2130, 0x213e, 0x2145, 0x2183, 0x2c00, 0x2c60, 0x2c62, 0x2c67, 0x2c69, 0x2c6b, 0x2c6d, 0x2c72, 0x2c75, 0x2c7e, + 0x2c82, 0x2c84, 0x2c86, 0x2c88, 0x2c8a, 0x2c8c, 0x2c8e, 0x2c90, 0x2c92, 0x2c94, 0x2c96, 0x2c98, 0x2c9a, 0x2c9c, 0x2c9e, 0x2ca0, 0x2ca2, + 0x2ca4, 0x2ca6, 0x2ca8, 0x2caa, 0x2cac, 0x2cae, 0x2cb0, 0x2cb2, 0x2cb4, 0x2cb6, 0x2cb8, 0x2cba, 0x2cbc, 0x2cbe, 0x2cc0, 0x2cc2, 0x2cc4, + 0x2cc6, 0x2cc8, 0x2cca, 0x2ccc, 0x2cce, 0x2cd0, 0x2cd2, 0x2cd4, 0x2cd6, 0x2cd8, 0x2cda, 0x2cdc, 0x2cde, 0x2ce0, 0x2ce2, 0x2ceb, 0x2ced, + 0x2cf2, 0xa640, 0xa642, 0xa644, 0xa646, 0xa648, 0xa64a, 0xa64c, 0xa64e, 0xa650, 0xa652, 0xa654, 0xa656, 0xa658, 0xa65a, 0xa65c, 0xa65e, + 0xa660, 0xa662, 0xa664, 0xa666, 0xa668, 0xa66a, 0xa66c, 0xa680, 0xa682, 0xa684, 0xa686, 0xa688, 0xa68a, 0xa68c, 0xa68e, 0xa690, 0xa692, + 0xa694, 0xa696, 0xa698, 0xa69a, 0xa722, 0xa724, 0xa726, 0xa728, 0xa72a, 0xa72c, 0xa72e, 0xa732, 0xa734, 0xa736, 0xa738, 0xa73a, 0xa73c, + 0xa73e, 0xa740, 0xa742, 0xa744, 0xa746, 0xa748, 0xa74a, 0xa74c, 0xa74e, 0xa750, 0xa752, 0xa754, 0xa756, 0xa758, 0xa75a, 0xa75c, 0xa75e, + 0xa760, 0xa762, 0xa764, 0xa766, 0xa768, 0xa76a, 0xa76c, 0xa76e, 0xa779, 0xa77b, 0xa77d, 0xa780, 0xa782, 0xa784, 0xa786, 0xa78b, 0xa78d, + 0xa790, 0xa792, 0xa796, 0xa798, 0xa79a, 0xa79c, 0xa79e, 0xa7a0, 0xa7a2, 0xa7a4, 0xa7a6, 0xa7a8, 0xa7aa, 0xa7b0, 0xa7b6, 0xa7b8, 0xa7ba, + 0xa7bc, 0xa7be, 0xa7c0, 0xa7c2, 0xa7c4, 0xa7c9, 0xa7cb, 0xa7d0, 0xa7d6, 0xa7d8, 0xa7da, 0xa7dc, 0xa7f5, 0xff21, + 0x10400, 0x104b0, 0x10570, 0x1057c, 0x1058c, 0x10594, 0x10c80, 0x10d50, 0x118a0, 0x16e40, 0x1d400, 0x1d434, 0x1d468, 0x1d49c, 0x1d49e, + 0x1d4a2, 0x1d4a5, 0x1d4a9, 0x1d4ae, 0x1d4d0, 0x1d504, 0x1d507, 0x1d50d, 0x1d516, 0x1d538, 0x1d53b, 0x1d540, 0x1d546, 0x1d54a, 0x1d56c, + 0x1d5a0, 0x1d5d4, 0x1d608, 0x1d63c, 0x1d670, 0x1d6a8, 0x1d6e2, 0x1d71c, 0x1d756, 0x1d790, 0x1d7ca, 0x1e900, +) + +private val upperCaseRangeEnds = intArrayOf( + 0x005a, 0x00d6, 0x00de, 0x0100, 0x0102, 0x0104, 0x0106, 0x0108, 0x010a, 0x010c, 0x010e, 0x0110, 0x0112, 0x0114, 0x0116, 0x0118, 0x011a, + 0x011c, 0x011e, 0x0120, 0x0122, 0x0124, 0x0126, 0x0128, 0x012a, 0x012c, 0x012e, 0x0130, 0x0132, 0x0134, 0x0136, 0x0139, 0x013b, 0x013d, + 0x013f, 0x0141, 0x0143, 0x0145, 0x0147, 0x014a, 0x014c, 0x014e, 0x0150, 0x0152, 0x0154, 0x0156, 0x0158, 0x015a, 0x015c, 0x015e, 0x0160, + 0x0162, 0x0164, 0x0166, 0x0168, 0x016a, 0x016c, 0x016e, 0x0170, 0x0172, 0x0174, 0x0176, 0x0179, 0x017b, 0x017d, 0x0182, 0x0184, 0x0187, + 0x018b, 0x0191, 0x0194, 0x0198, 0x019d, 0x01a0, 0x01a2, 0x01a4, 0x01a7, 0x01a9, 0x01ac, 0x01af, 0x01b3, 0x01b5, 0x01b8, 0x01bc, 0x01c4, + 0x01c7, 0x01ca, 0x01cd, 0x01cf, 0x01d1, 0x01d3, 0x01d5, 0x01d7, 0x01d9, 0x01db, 0x01de, 0x01e0, 0x01e2, 0x01e4, 0x01e6, 0x01e8, 0x01ea, + 0x01ec, 0x01ee, 0x01f1, 0x01f4, 0x01f8, 0x01fa, 0x01fc, 0x01fe, 0x0200, 0x0202, 0x0204, 0x0206, 0x0208, 0x020a, 0x020c, 0x020e, 0x0210, + 0x0212, 0x0214, 0x0216, 0x0218, 0x021a, 0x021c, 0x021e, 0x0220, 0x0222, 0x0224, 0x0226, 0x0228, 0x022a, 0x022c, 0x022e, 0x0230, 0x0232, + 0x023b, 0x023e, 0x0241, 0x0246, 0x0248, 0x024a, 0x024c, 0x024e, 0x0370, 0x0372, 0x0376, 0x037f, 0x0386, 0x038a, 0x038c, 0x038f, 0x03a1, + 0x03ab, 0x03cf, 0x03d4, 0x03d8, 0x03da, 0x03dc, 0x03de, 0x03e0, 0x03e2, 0x03e4, 0x03e6, 0x03e8, 0x03ea, 0x03ec, 0x03ee, 0x03f4, 0x03f7, + 0x03fa, 0x042f, 0x0460, 0x0462, 0x0464, 0x0466, 0x0468, 0x046a, 0x046c, 0x046e, 0x0470, 0x0472, 0x0474, 0x0476, 0x0478, 0x047a, 0x047c, + 0x047e, 0x0480, 0x048a, 0x048c, 0x048e, 0x0490, 0x0492, 0x0494, 0x0496, 0x0498, 0x049a, 0x049c, 0x049e, 0x04a0, 0x04a2, 0x04a4, 0x04a6, + 0x04a8, 0x04aa, 0x04ac, 0x04ae, 0x04b0, 0x04b2, 0x04b4, 0x04b6, 0x04b8, 0x04ba, 0x04bc, 0x04be, 0x04c1, 0x04c3, 0x04c5, 0x04c7, 0x04c9, + 0x04cb, 0x04cd, 0x04d0, 0x04d2, 0x04d4, 0x04d6, 0x04d8, 0x04da, 0x04dc, 0x04de, 0x04e0, 0x04e2, 0x04e4, 0x04e6, 0x04e8, 0x04ea, 0x04ec, + 0x04ee, 0x04f0, 0x04f2, 0x04f4, 0x04f6, 0x04f8, 0x04fa, 0x04fc, 0x04fe, 0x0500, 0x0502, 0x0504, 0x0506, 0x0508, 0x050a, 0x050c, 0x050e, + 0x0510, 0x0512, 0x0514, 0x0516, 0x0518, 0x051a, 0x051c, 0x051e, 0x0520, 0x0522, 0x0524, 0x0526, 0x0528, 0x052a, 0x052c, 0x052e, 0x0556, + 0x10c5, 0x10c7, 0x10cd, 0x13f5, 0x1c89, 0x1cba, 0x1cbf, 0x1e00, 0x1e02, 0x1e04, 0x1e06, 0x1e08, 0x1e0a, 0x1e0c, 0x1e0e, 0x1e10, 0x1e12, + 0x1e14, 0x1e16, 0x1e18, 0x1e1a, 0x1e1c, 0x1e1e, 0x1e20, 0x1e22, 0x1e24, 0x1e26, 0x1e28, 0x1e2a, 0x1e2c, 0x1e2e, 0x1e30, 0x1e32, 0x1e34, + 0x1e36, 0x1e38, 0x1e3a, 0x1e3c, 0x1e3e, 0x1e40, 0x1e42, 0x1e44, 0x1e46, 0x1e48, 0x1e4a, 0x1e4c, 0x1e4e, 0x1e50, 0x1e52, 0x1e54, 0x1e56, + 0x1e58, 0x1e5a, 0x1e5c, 0x1e5e, 0x1e60, 0x1e62, 0x1e64, 0x1e66, 0x1e68, 0x1e6a, 0x1e6c, 0x1e6e, 0x1e70, 0x1e72, 0x1e74, 0x1e76, 0x1e78, + 0x1e7a, 0x1e7c, 0x1e7e, 0x1e80, 0x1e82, 0x1e84, 0x1e86, 0x1e88, 0x1e8a, 0x1e8c, 0x1e8e, 0x1e90, 0x1e92, 0x1e94, 0x1e9e, 0x1ea0, 0x1ea2, + 0x1ea4, 0x1ea6, 0x1ea8, 0x1eaa, 0x1eac, 0x1eae, 0x1eb0, 0x1eb2, 0x1eb4, 0x1eb6, 0x1eb8, 0x1eba, 0x1ebc, 0x1ebe, 0x1ec0, 0x1ec2, 0x1ec4, + 0x1ec6, 0x1ec8, 0x1eca, 0x1ecc, 0x1ece, 0x1ed0, 0x1ed2, 0x1ed4, 0x1ed6, 0x1ed8, 0x1eda, 0x1edc, 0x1ede, 0x1ee0, 0x1ee2, 0x1ee4, 0x1ee6, + 0x1ee8, 0x1eea, 0x1eec, 0x1eee, 0x1ef0, 0x1ef2, 0x1ef4, 0x1ef6, 0x1ef8, 0x1efa, 0x1efc, 0x1efe, 0x1f0f, 0x1f1d, 0x1f2f, 0x1f3f, 0x1f4d, + 0x1f59, 0x1f5b, 0x1f5d, 0x1f5f, 0x1f6f, 0x1fbb, 0x1fcb, 0x1fdb, 0x1fec, 0x1ffb, 0x2102, 0x2107, 0x210d, 0x2112, 0x2115, 0x211d, 0x2124, + 0x2126, 0x2128, 0x212d, 0x2133, 0x213f, 0x2145, 0x2183, 0x2c2f, 0x2c60, 0x2c64, 0x2c67, 0x2c69, 0x2c6b, 0x2c70, 0x2c72, 0x2c75, 0x2c80, + 0x2c82, 0x2c84, 0x2c86, 0x2c88, 0x2c8a, 0x2c8c, 0x2c8e, 0x2c90, 0x2c92, 0x2c94, 0x2c96, 0x2c98, 0x2c9a, 0x2c9c, 0x2c9e, 0x2ca0, 0x2ca2, + 0x2ca4, 0x2ca6, 0x2ca8, 0x2caa, 0x2cac, 0x2cae, 0x2cb0, 0x2cb2, 0x2cb4, 0x2cb6, 0x2cb8, 0x2cba, 0x2cbc, 0x2cbe, 0x2cc0, 0x2cc2, 0x2cc4, + 0x2cc6, 0x2cc8, 0x2cca, 0x2ccc, 0x2cce, 0x2cd0, 0x2cd2, 0x2cd4, 0x2cd6, 0x2cd8, 0x2cda, 0x2cdc, 0x2cde, 0x2ce0, 0x2ce2, 0x2ceb, 0x2ced, + 0x2cf2, 0xa640, 0xa642, 0xa644, 0xa646, 0xa648, 0xa64a, 0xa64c, 0xa64e, 0xa650, 0xa652, 0xa654, 0xa656, 0xa658, 0xa65a, 0xa65c, 0xa65e, + 0xa660, 0xa662, 0xa664, 0xa666, 0xa668, 0xa66a, 0xa66c, 0xa680, 0xa682, 0xa684, 0xa686, 0xa688, 0xa68a, 0xa68c, 0xa68e, 0xa690, 0xa692, + 0xa694, 0xa696, 0xa698, 0xa69a, 0xa722, 0xa724, 0xa726, 0xa728, 0xa72a, 0xa72c, 0xa72e, 0xa732, 0xa734, 0xa736, 0xa738, 0xa73a, 0xa73c, + 0xa73e, 0xa740, 0xa742, 0xa744, 0xa746, 0xa748, 0xa74a, 0xa74c, 0xa74e, 0xa750, 0xa752, 0xa754, 0xa756, 0xa758, 0xa75a, 0xa75c, 0xa75e, + 0xa760, 0xa762, 0xa764, 0xa766, 0xa768, 0xa76a, 0xa76c, 0xa76e, 0xa779, 0xa77b, 0xa77e, 0xa780, 0xa782, 0xa784, 0xa786, 0xa78b, 0xa78d, + 0xa790, 0xa792, 0xa796, 0xa798, 0xa79a, 0xa79c, 0xa79e, 0xa7a0, 0xa7a2, 0xa7a4, 0xa7a6, 0xa7a8, 0xa7ae, 0xa7b4, 0xa7b6, 0xa7b8, 0xa7ba, + 0xa7bc, 0xa7be, 0xa7c0, 0xa7c2, 0xa7c7, 0xa7c9, 0xa7cc, 0xa7d0, 0xa7d6, 0xa7d8, 0xa7da, 0xa7dc, 0xa7f5, 0xff3a, + 0x10427, 0x104d3, 0x1057a, 0x1058a, 0x10592, 0x10595, 0x10cb2, 0x10d65, 0x118bf, 0x16e5f, 0x1d419, 0x1d44d, 0x1d481, 0x1d49c, 0x1d49f, + 0x1d4a2, 0x1d4a6, 0x1d4ac, 0x1d4b5, 0x1d4e9, 0x1d505, 0x1d50a, 0x1d514, 0x1d51c, 0x1d539, 0x1d53e, 0x1d544, 0x1d546, 0x1d550, 0x1d585, + 0x1d5b9, 0x1d5ed, 0x1d621, 0x1d655, 0x1d689, 0x1d6c0, 0x1d6fa, 0x1d734, 0x1d76e, 0x1d7a8, 0x1d7ca, 0x1e921, +) /* auxiliary function for binary search in interval table */ private fun bisearch(ucs: Char, table: Array, max: Int): Int { diff --git a/fleet/kernel/src/fleet/kernel/ReteExt.kt b/fleet/kernel/src/fleet/kernel/ReteExt.kt index c8af644a8043..f85c2440a4bb 100644 --- a/fleet/kernel/src/fleet/kernel/ReteExt.kt +++ b/fleet/kernel/src/fleet/kernel/ReteExt.kt @@ -98,7 +98,7 @@ suspend fun waitFor(p: () -> Boolean) { val result = queryAsFlow { p() }.firstOrNull { it } // query could be terminated before our coroutine, null means we are in shutdown if (result == null) { - throw CancellationException() + throw CancellationException("Query was terminated") } } @@ -108,7 +108,7 @@ private object Logger { suspend fun waitForNotNullWithTimeout(timeMillis: Long = 30000L, p: () -> T?): T { return waitForNotNullWithTimeoutOrNull(timeMillis, p) ?: run { - Logger.logger.error(Throwable("Timed out waiting for ${p.javaClass} to return not null, $timeMillis ms")) + Logger.logger.error(Throwable("Timed out waiting for ${p::class} to return not null, $timeMillis ms")) throw CancellationException("$p is null, after ${timeMillis}ms") } } diff --git a/fleet/kernel/src/fleet/kernel/rebase/RebaseLoop.kt b/fleet/kernel/src/fleet/kernel/rebase/RebaseLoop.kt index b888a4779d1b..313fc3e65643 100644 --- a/fleet/kernel/src/fleet/kernel/rebase/RebaseLoop.kt +++ b/fleet/kernel/src/fleet/kernel/rebase/RebaseLoop.kt @@ -228,7 +228,7 @@ private fun ChangeScope.runEffects(list: List) { e.effect(this) } catch (x: Throwable) { - logger.error(x, "failed running effect ${e.javaClass} in offer") + logger.error(x, "failed running effect ${e::class} in offer") } } } diff --git a/fleet/kernel/src/fleet/kernel/rete/WithMatchResult.kt b/fleet/kernel/src/fleet/kernel/rete/WithMatchResult.kt index bdfa43f5b8ff..a5c187d5c561 100644 --- a/fleet/kernel/src/fleet/kernel/rete/WithMatchResult.kt +++ b/fleet/kernel/src/fleet/kernel/rete/WithMatchResult.kt @@ -1,8 +1,8 @@ // Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package fleet.kernel.rete +import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CopyableThrowable -import kotlin.coroutines.cancellation.CancellationException sealed class WithMatchResult { data class Success(val value: T) : WithMatchResult() @@ -38,7 +38,11 @@ class CancellationReason(val reason: String, val match: Match<*>?) class UnsatisfiedMatchException(val reason: CancellationReason) : CancellationException(reason.reason), CopyableThrowable { + + override var cause: Throwable? = super.cause + private set + override fun createCopy(): UnsatisfiedMatchException { - return UnsatisfiedMatchException(reason).also { it.initCause(this) } + return UnsatisfiedMatchException(reason).also { it.cause = this } } } \ No newline at end of file diff --git a/fleet/kernel/src/fleet/kernel/rete/impl/Tracing.kt b/fleet/kernel/src/fleet/kernel/rete/impl/Tracing.kt index 6174ea9736be..6ef01b41b267 100644 --- a/fleet/kernel/src/fleet/kernel/rete/impl/Tracing.kt +++ b/fleet/kernel/src/fleet/kernel/rete/impl/Tracing.kt @@ -33,7 +33,7 @@ internal fun Query.tracing(trackingKey: QueryTracingKey?): Query = let } val producer = query.producer() Producer { emit -> - val collectorId = System.identityHashCode(emit) + val collectorId = emit.hashCode() trackingKey.logger.info { "${trackingKey}: query $queryId is being collected by $collectorId" } onDispose { trackingKey.logger.info { "${trackingKey}: query $queryId has stopped being collected by $collectorId" } diff --git a/fleet/rpc/src/fleet/rpc/client/RpcClient.kt b/fleet/rpc/src/fleet/rpc/client/RpcClient.kt index be8b667e7d39..4435cb10cab0 100644 --- a/fleet/rpc/src/fleet/rpc/client/RpcClient.kt +++ b/fleet/rpc/src/fleet/rpc/client/RpcClient.kt @@ -26,7 +26,6 @@ import kotlinx.coroutines.flow.* import kotlinx.coroutines.selects.whileSelect import kotlinx.serialization.builtins.serializer import kotlinx.serialization.json.Json -import java.util.Optional import java.util.concurrent.ConcurrentHashMap import kotlin.coroutines.* @@ -368,7 +367,7 @@ class RpcClient internal constructor( it.span.addEvent("failure").end() } } - streams.values.removeIf { + streams.values.removeAll { it.closeStream(throwable) true } @@ -385,7 +384,7 @@ class RpcClient internal constructor( } } } - streams.values.removeIf { + streams.values.removeAll { if (it.route == route) { it.closeStream(RouteClosedException(route, rpcStreamFailureMessage(it.displayName, message))) true @@ -478,7 +477,7 @@ class RpcClient internal constructor( } is InternalStreamDescriptor.ToRemote -> { val cancellationException = (cause as? CancellationException) ?: cause?.let { CancellationException(it.message, it) } - budget.cancel(cancellationException ?: CancellationException()) + budget.cancel(cancellationException ?: CancellationException("The stream was cancelled")) channel.cancel(cancellationException) } } @@ -516,98 +515,96 @@ class RpcClient internal constructor( service = call.service, method = call.signature.methodName, args = serializedArguments) - withTimeoutOrNull(RPC_TIMEOUT) { - optional { - val callRequest = requestInterceptor.interceptCallRequest(uninterceptedRequest) - logger.trace { "Interceptor completed for request ${callRequest}" } - val rpcStrategy = coroutineContext[RpcStrategyContextElement] ?: RpcStrategyContextElement() - if (rpcStrategy.awaitConnection) { - logger.trace { "request $requestId, waiting for ${call.route} to become available" } - grayList[call.route]?.await() - logger.trace { "request $requestId, ${call.route} is available" } - } - else if (grayList.contains(call.route)) { - throw RouteClosedException(call.route, rpcCallFailureMessage(callRequest, "Route ${call.route} closed")) - } + withTimeoutOrNull(RPC_TIMEOUT) { + val callRequest = requestInterceptor.interceptCallRequest(uninterceptedRequest) + logger.trace { "Interceptor completed for request ${callRequest}" } + val rpcStrategy = coroutineContext[RpcStrategyContextElement] ?: RpcStrategyContextElement() + if (rpcStrategy.awaitConnection) { + logger.trace { "request $requestId, waiting for ${call.route} to become available" } + grayList[call.route]?.await() + logger.trace { "request $requestId, ${call.route} is available" } + } + else if (grayList.contains(call.route)) { + throw RouteClosedException(call.route, rpcCallFailureMessage(callRequest, "Route ${call.route} closed")) + } - val span = tracer.spanBuilder("RPC Call") - .setAttribute("service", callRequest.service.id) - .setAttribute("method", callRequest.method) - .setSpanKind(SpanKind.CLIENT) - .startSpan() - val otelData = Context.current().with(span).toTelemetryData() - suspendCancellableCoroutine { cc -> - val request = OutgoingRequest(route = call.route, - call = callRequest, - token = token, - continuation = cc, - returnType = call.signature.returnType, - streamParameters = streamParameters, - prefetchStrategy = rpcStrategy.prefetchStrategy) - val resumeWithException = { cause: Throwable -> - val exToResumeWith = cause.causeOfType()?.let { RpcClientDisconnectedException(null, it) } - ?: cause - logger.trace(exToResumeWith) { "Failed to send request $requestId with exception, remove it from queue" } - outgoingRpc.remove(requestId)?.let { (r, span) -> - span.setStatus(StatusCode.ERROR, "Failed to send call request: ${exToResumeWith.stackTraceToString()}").end() - for (stream in r.streamParameters) { - unregisterStream(stream.uid) { exToResumeWith } - } - r.continuation.resumeWithException(exToResumeWith) + val span = tracer.spanBuilder("RPC Call") + .setAttribute("service", callRequest.service.id) + .setAttribute("method", callRequest.method) + .setSpanKind(SpanKind.CLIENT) + .startSpan() + val otelData = Context.current().with(span).toTelemetryData() + suspendCancellableCoroutine { cc -> + val request = OutgoingRequest(route = call.route, + call = callRequest, + token = token, + continuation = cc, + returnType = call.signature.returnType, + streamParameters = streamParameters, + prefetchStrategy = rpcStrategy.prefetchStrategy) + val resumeWithException = { cause: Throwable -> + val exToResumeWith = cause.causeOfType()?.let { RpcClientDisconnectedException(null, it) } + ?: cause + logger.trace(exToResumeWith) { "Failed to send request $requestId with exception, remove it from queue" } + outgoingRpc.remove(requestId)?.let { (r, span) -> + span.setStatus(StatusCode.ERROR, "Failed to send call request: ${exToResumeWith.stackTraceToString()}").end() + for (stream in r.streamParameters) { + unregisterStream(stream.uid) { exToResumeWith } } + r.continuation.resumeWithException(exToResumeWith) } - executeCommand { cause -> - if (cause == null) { - logger.trace { "Register request ${request.call} in queue" } - val previous = outgoingRpc.putIfAbsent(requestId, OngoingRequest(request, span)) - check(previous == null) { "Request with id $requestId is already present in the queue" } - val streamDescriptors = registerStreams(request.streamParameters, request.route, rpcStrategy.prefetchStrategy) - sendAsync(callRequest.seal(destination = request.route, origin = origin, otelData = otelData)) { cause -> - if (cause == null) { - logger.trace { "Request sent ${request.call}" } - // register cancellation handler only after the request is enqueued or we can end up sending CancelCall before CallRequest - request.continuation.invokeOnCancellation { c -> - if (c != null) { - // be careful, invokeOnCancellation is invoked concurrently to the main event loop - executeCommand { ex -> - if (ex == null) { - requestCanceledByClient(requestId, c) - } + } + executeCommand { cause -> + if (cause == null) { + logger.trace { "Register request ${request.call} in queue" } + val previous = outgoingRpc.putIfAbsent(requestId, OngoingRequest(request, span)) + check(previous == null) { "Request with id $requestId is already present in the queue" } + val streamDescriptors = registerStreams(request.streamParameters, request.route, rpcStrategy.prefetchStrategy) + sendAsync(callRequest.seal(destination = request.route, origin = origin, otelData = otelData)) { cause -> + if (cause == null) { + logger.trace { "Request sent ${request.call}" } + // register cancellation handler only after the request is enqueued or we can end up sending CancelCall before CallRequest + request.continuation.invokeOnCancellation { c -> + if (c != null) { + // be careful, invokeOnCancellation is invoked concurrently to the main event loop + executeCommand { ex -> + if (ex == null) { + requestCanceledByClient(requestId, c) } } } - for (internalDescriptor in streamDescriptors) { - serveStream(internalDescriptor, rpcStrategy.prefetchStrategy) - } } - else { - resumeWithException(cause) + for (internalDescriptor in streamDescriptors) { + serveStream(internalDescriptor, rpcStrategy.prefetchStrategy) } } - } - else { - val exToResumeWith = cause.causeOfType()?.let { RpcClientDisconnectedException(null, it) } - ?: cause - cc.resumeWithException(exToResumeWith) + else { + resumeWithException(cause) + } } } - }.let { result -> - logger.trace { "Resumed $requestId, serving response streams" } - // resumed successfully, start serving streams - val resource = completedRpc.remove(requestId)?.also { - for (internalDescriptor in it.streams) { - serveStream(internalDescriptor, it.prefetchStrategy) - } - } ?: run { - logger.trace { "No resources assigned for $requestId, was it cancelled already?" } - null + else { + val exToResumeWith = cause.causeOfType()?.let { RpcClientDisconnectedException(null, it) } + ?: cause + cc.resumeWithException(exToResumeWith) } - val disposable = SuspendInvocationHandler.CallResult(result) { - resource?.let(::disposeResponseResource) - } - publish(disposable) - logger.trace { "Result published for request $requestId" } } + }.let { result -> + logger.trace { "Resumed $requestId, serving response streams" } + // resumed successfully, start serving streams + val resource = completedRpc.remove(requestId)?.also { + for (internalDescriptor in it.streams) { + serveStream(internalDescriptor, it.prefetchStrategy) + } + } ?: run { + logger.trace { "No resources assigned for $requestId, was it cancelled already?" } + null + } + val disposable = SuspendInvocationHandler.CallResult(result) { + resource?.let(::disposeResponseResource) + } + publish(disposable) + logger.trace { "Result published for request $requestId" } } } ?: throw RpcTimeoutException("Request $uninterceptedRequest has timed out after ${RPC_TIMEOUT}ms", cause = null) } @@ -640,7 +637,3 @@ class RpcClient internal constructor( sendAsync = ::sendAsync) } } - -internal inline fun optional(body: () -> T?): Optional { - return Optional.ofNullable(body()) -} \ No newline at end of file diff --git a/fleet/rpc/src/fleet/rpc/core/RpcStream.kt b/fleet/rpc/src/fleet/rpc/core/RpcStream.kt index 33f9560b9721..27a2bc58f891 100644 --- a/fleet/rpc/src/fleet/rpc/core/RpcStream.kt +++ b/fleet/rpc/src/fleet/rpc/core/RpcStream.kt @@ -244,7 +244,10 @@ fun serveStream( // register streams before we publish them to remote with `sendAsync` or we may miss some messages from FROM_REMOTE streams if remote is fast enough val internalStreamDescriptors = streamDescriptors.map { registerStream(it) } descriptor.budget.withdrawSuspend() - RpcStream.logger.trace { "Sending in stream ${descriptor.uid} <${descriptor.displayName}> item ${item?.javaClass?.simpleName}($item)" } + RpcStream.logger.trace { + val itemName = item?.let { it::class.simpleName } + "Sending in stream ${descriptor.uid} <${descriptor.displayName}> item ${itemName}($item)" + } sendMessage(RpcMessage.StreamData(streamId = descriptor.uid, data = jsonElement)) // we must serve TO_REMOTE streams only after initial message was sent for (internalStream in internalStreamDescriptors) {