| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 11c9b0e commit ca76cb7
6 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1,4 +1,8 @@ | |||
| 1 | + import _collections_abc | ||
| 1 | 2 | from _contextvars import Context, ContextVar, Token, copy_context | |
| 2 | 3 | ||
| 3 | 4 | ||
| 4 | 5 | __all__ = ('Context', 'ContextVar', 'Token', 'copy_context') | |
| 6 | + | ||
| 7 | + | ||
| 8 | + _collections_abc.Mapping.register(Context) | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -889,10 +889,11 @@ def test_sys_flags(self): | |||
| 889 | 889 | "dont_write_bytecode", "no_user_site", "no_site", | |
| 890 | 890 | "ignore_environment", "verbose", "bytes_warning", "quiet", | |
| 891 | 891 | "hash_randomization", "isolated", "dev_mode", "utf8_mode", | |
| 892 | - "warn_default_encoding", "safe_path", "int_max_str_digits") | ||
| 892 | + "warn_default_encoding", "safe_path", "int_max_str_digits", | ||
| 893 | + "thread_inherit_context") | ||
| 893 | 894 | for attr in attrs: | |
| 894 | 895 | self.assertTrue(hasattr(sys.flags, attr), attr) | |
| 895 | - attr_type = bool if attr in ("dev_mode", "safe_path") else int | ||
| 896 | + attr_type = bool if attr in ("dev_mode", "safe_path", "thread_inherit_context") else int | ||
| 896 | 897 | self.assertEqual(type(getattr(sys.flags, attr)), attr_type, attr) | |
| 897 | 898 | self.assertTrue(repr(sys.flags)) | |
| 898 | 899 | self.assertEqual(len(sys.flags), len(attrs)) | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -3,7 +3,7 @@ | |||
| 3 | 3 | import os as _os | |
| 4 | 4 | import sys as _sys | |
| 5 | 5 | import _thread | |
| 6 | - import warnings | ||
| 6 | + import _contextvars | ||
| 7 | 7 | ||
| 8 | 8 | from time import monotonic as _time | |
| 9 | 9 | from _weakrefset import WeakSet | |
@@ -48,6 +48,10 @@ | |||
| 48 | 48 | __all__.append('get_native_id') | |
| 49 | 49 | except AttributeError: | |
| 50 | 50 | _HAVE_THREAD_NATIVE_ID = False | |
| 51 | + try: | ||
| 52 | + _set_name = _thread.set_name | ||
| 53 | + except AttributeError: | ||
| 54 | + _set_name = None | ||
| 51 | 55 | ThreadError = _thread.error | |
| 52 | 56 | try: | |
| 53 | 57 | _CRLock = _thread.RLock | |
@@ -129,6 +133,7 @@ def RLock(*args, **kwargs): | |||
| 129 | 133 | ||
| 130 | 134 | """ | |
| 131 | 135 | if args or kwargs: | |
| 136 | + import warnings | ||
| 132 | 137 | warnings.warn( | |
| 133 | 138 | 'Passing arguments to RLock is deprecated and will be removed in 3.15', | |
| 134 | 139 | DeprecationWarning, | |
@@ -160,7 +165,7 @@ def __repr__(self): | |||
| 160 | 165 | except KeyError: | |
| 161 | 166 | pass | |
| 162 | 167 | return "<%s %s.%s object owner=%r count=%d at %s>" % ( | |
| 163 | - "locked" if self._block.locked() else "unlocked", | ||
| 168 | + "locked" if self.locked() else "unlocked", | ||
| 164 | 169 | self.__class__.__module__, | |
| 165 | 170 | self.__class__.__qualname__, | |
| 166 | 171 | owner, | |
@@ -237,6 +242,10 @@ def release(self): | |||
| 237 | 242 | def __exit__(self, t, v, tb): | |
| 238 | 243 | self.release() | |
| 239 | 244 | ||
| 245 | + def locked(self): | ||
| 246 | + """Return whether this object is locked.""" | ||
| 247 | + return self._block.locked() | ||
| 248 | + | ||
| 240 | 249 | # Internal methods used by condition variables | |
| 241 | 250 | ||
| 242 | 251 | def _acquire_restore(self, state): | |
@@ -282,9 +291,10 @@ def __init__(self, lock=None): | |||
| 282 | 291 | if lock is None: | |
| 283 | 292 | lock = RLock() | |
| 284 | 293 | self._lock = lock | |
| 285 | - # Export the lock's acquire() and release() methods | ||
| 294 | + # Export the lock's acquire(), release(), and locked() methods | ||
| 286 | 295 | self.acquire = lock.acquire | |
| 287 | 296 | self.release = lock.release | |
| 297 | + self.locked = lock.locked | ||
| 288 | 298 | # If the lock defines _release_save() and/or _acquire_restore(), | |
| 289 | 299 | # these override the default implementations (which just call | |
| 290 | 300 | # release() and acquire() on the lock). Ditto for _is_owned(). | |
@@ -868,7 +878,7 @@ class Thread: | |||
| 868 | 878 | _initialized = False | |
| 869 | 879 | ||
| 870 | 880 | def __init__(self, group=None, target=None, name=None, | |
| 871 | - args=(), kwargs=None, *, daemon=None): | ||
| 881 | + args=(), kwargs=None, *, daemon=None, context=None): | ||
| 872 | 882 | """This constructor should always be called with keyword arguments. Arguments are: | |
| 873 | 883 | ||
| 874 | 884 | *group* should be None; reserved for future extension when a ThreadGroup | |
@@ -885,6 +895,14 @@ class is implemented. | |||
| 885 | 895 | *kwargs* is a dictionary of keyword arguments for the target | |
| 886 | 896 | invocation. Defaults to {}. | |
| 887 | 897 | ||
| 898 | + *context* is the contextvars.Context value to use for the thread. | ||
| 899 | + The default value is None, which means to check | ||
| 900 | + sys.flags.thread_inherit_context. If that flag is true, use a copy | ||
| 901 | + of the context of the caller. If false, use an empty context. To | ||
| 902 | + explicitly start with an empty context, pass a new instance of | ||
| 903 | + contextvars.Context(). To explicitly start with a copy of the current | ||
| 904 | + context, pass the value from contextvars.copy_context(). | ||
| 905 | + | ||
| 888 | 906 | If a subclass overrides the constructor, it must make sure to invoke | |
| 889 | 907 | the base class constructor (Thread.__init__()) before doing anything | |
| 890 | 908 | else to the thread. | |
@@ -914,10 +932,11 @@ class is implemented. | |||
| 914 | 932 | self._daemonic = daemon | |
| 915 | 933 | else: | |
| 916 | 934 | self._daemonic = current_thread().daemon | |
| 935 | + self._context = context | ||
| 917 | 936 | self._ident = None | |
| 918 | 937 | if _HAVE_THREAD_NATIVE_ID: | |
| 919 | 938 | self._native_id = None | |
| 920 | - self._handle = _ThreadHandle() | ||
| 939 | + self._os_thread_handle = _ThreadHandle() | ||
| 921 | 940 | self._started = Event() | |
| 922 | 941 | self._initialized = True | |
| 923 | 942 | # Copy of sys.stderr used by self._invoke_excepthook() | |
@@ -932,7 +951,7 @@ def _after_fork(self, new_ident=None): | |||
| 932 | 951 | if new_ident is not None: | |
| 933 | 952 | # This thread is alive. | |
| 934 | 953 | self._ident = new_ident | |
| 935 | - assert self._handle.ident == new_ident | ||
| 954 | + assert self._os_thread_handle.ident == new_ident | ||
| 936 | 955 | if _HAVE_THREAD_NATIVE_ID: | |
| 937 | 956 | self._set_native_id() | |
| 938 | 957 | else: | |
@@ -945,7 +964,7 @@ def __repr__(self): | |||
| 945 | 964 | status = "initial" | |
| 946 | 965 | if self._started.is_set(): | |
| 947 | 966 | status = "started" | |
| 948 | - if self._handle.is_done(): | ||
| 967 | + if self._os_thread_handle.is_done(): | ||
| 949 | 968 | status = "stopped" | |
| 950 | 969 | if self._daemonic: | |
| 951 | 970 | status += " daemon" | |
@@ -971,9 +990,19 @@ def start(self): | |||
| 971 | 990 | ||
| 972 | 991 | with _active_limbo_lock: | |
| 973 | 992 | _limbo[self] = self | |
| 993 | + | ||
| 994 | + if self._context is None: | ||
| 995 | + # No context provided | ||
| 996 | + if _sys.flags.thread_inherit_context: | ||
| 997 | + # start with a copy of the context of the caller | ||
| 998 | + self._context = _contextvars.copy_context() | ||
| 999 | + else: | ||
| 1000 | + # start with an empty context | ||
| 1001 | + self._context = _contextvars.Context() | ||
| 1002 | + | ||
| 974 | 1003 | try: | |
| 975 | 1004 | # Start joinable thread | |
| 976 | - _start_joinable_thread(self._bootstrap, handle=self._handle, | ||
| 1005 | + _start_joinable_thread(self._bootstrap, handle=self._os_thread_handle, | ||
| 977 | 1006 | daemon=self.daemon) | |
| 978 | 1007 | except Exception: | |
| 979 | 1008 | with _active_limbo_lock: | |
@@ -1025,11 +1054,20 @@ def _set_ident(self): | |||
| 1025 | 1054 | def _set_native_id(self): | |
| 1026 | 1055 | self._native_id = get_native_id() | |
| 1027 | 1056 | ||
| 1057 | + def _set_os_name(self): | ||
| 1058 | + if _set_name is None or not self._name: | ||
| 1059 | + return | ||
| 1060 | + try: | ||
| 1061 | + _set_name(self._name) | ||
| 1062 | + except OSError: | ||
| 1063 | + pass | ||
| 1064 | + | ||
| 1028 | 1065 | def _bootstrap_inner(self): | |
| 1029 | 1066 | try: | |
| 1030 | 1067 | self._set_ident() | |
| 1031 | 1068 | if _HAVE_THREAD_NATIVE_ID: | |
| 1032 | 1069 | self._set_native_id() | |
| 1070 | + self._set_os_name() | ||
| 1033 | 1071 | self._started.set() | |
| 1034 | 1072 | with _active_limbo_lock: | |
| 1035 | 1073 | _active[self._ident] = self | |
@@ -1041,7 +1079,7 @@ def _bootstrap_inner(self): | |||
| 1041 | 1079 | _sys.setprofile(_profile_hook) | |
| 1042 | 1080 | ||
| 1043 | 1081 | try: | |
| 1044 | - self.run() | ||
| 1082 | + self._context.run(self.run) | ||
| 1045 | 1083 | except: | |
| 1046 | 1084 | self._invoke_excepthook(self) | |
| 1047 | 1085 | finally: | |
@@ -1092,7 +1130,7 @@ def join(self, timeout=None): | |||
| 1092 | 1130 | if timeout is not None: | |
| 1093 | 1131 | timeout = max(timeout, 0) | |
| 1094 | 1132 | ||
| 1095 | - self._handle.join(timeout) | ||
| 1133 | + self._os_thread_handle.join(timeout) | ||
| 1096 | 1134 | ||
| 1097 | 1135 | @property | |
| 1098 | 1136 | def name(self): | |
@@ -1109,6 +1147,8 @@ def name(self): | |||
| 1109 | 1147 | def name(self, name): | |
| 1110 | 1148 | assert self._initialized, "Thread.__init__() not called" | |
| 1111 | 1149 | self._name = str(name) | |
| 1150 | + if get_ident() == self._ident: | ||
| 1151 | + self._set_os_name() | ||
| 1112 | 1152 | ||
| 1113 | 1153 | @property | |
| 1114 | 1154 | def ident(self): | |
@@ -1143,7 +1183,7 @@ def is_alive(self): | |||
| 1143 | 1183 | ||
| 1144 | 1184 | """ | |
| 1145 | 1185 | assert self._initialized, "Thread.__init__() not called" | |
| 1146 | - return self._started.is_set() and not self._handle.is_done() | ||
| 1186 | + return self._started.is_set() and not self._os_thread_handle.is_done() | ||
| 1147 | 1187 | ||
| 1148 | 1188 | @property | |
| 1149 | 1189 | def daemon(self): | |
@@ -1354,7 +1394,7 @@ def __init__(self): | |||
| 1354 | 1394 | Thread.__init__(self, name="MainThread", daemon=False) | |
| 1355 | 1395 | self._started.set() | |
| 1356 | 1396 | self._ident = _get_main_thread_ident() | |
| 1357 | - self._handle = _make_thread_handle(self._ident) | ||
| 1397 | + self._os_thread_handle = _make_thread_handle(self._ident) | ||
| 1358 | 1398 | if _HAVE_THREAD_NATIVE_ID: | |
| 1359 | 1399 | self._set_native_id() | |
| 1360 | 1400 | with _active_limbo_lock: | |
@@ -1402,15 +1442,15 @@ def __init__(self): | |||
| 1402 | 1442 | daemon=_daemon_threads_allowed()) | |
| 1403 | 1443 | self._started.set() | |
| 1404 | 1444 | self._set_ident() | |
| 1405 | - self._handle = _make_thread_handle(self._ident) | ||
| 1445 | + self._os_thread_handle = _make_thread_handle(self._ident) | ||
| 1406 | 1446 | if _HAVE_THREAD_NATIVE_ID: | |
| 1407 | 1447 | self._set_native_id() | |
| 1408 | 1448 | with _active_limbo_lock: | |
| 1409 | 1449 | _active[self._ident] = self | |
| 1410 | 1450 | _DeleteDummyThreadOnDel(self) | |
| 1411 | 1451 | ||
| 1412 | 1452 | def is_alive(self): | |
| 1413 | - if not self._handle.is_done() and self._started.is_set(): | ||
| 1453 | + if not self._os_thread_handle.is_done() and self._started.is_set(): | ||
| 1414 | 1454 | return True | |
| 1415 | 1455 | raise RuntimeError("thread is not alive") | |
| 1416 | 1456 | ||
@@ -1524,7 +1564,7 @@ def _shutdown(): | |||
| 1524 | 1564 | # dubious, but some code does it. We can't wait for it to be marked as done | |
| 1525 | 1565 | # normally - that won't happen until the interpreter is nearly dead. So | |
| 1526 | 1566 | # mark it done here. | |
| 1527 | - if _main_thread._handle.is_done() and _is_main_interpreter(): | ||
| 1567 | + if _main_thread._os_thread_handle.is_done() and _is_main_interpreter(): | ||
| 1528 | 1568 | # _shutdown() was already called | |
| 1529 | 1569 | return | |
| 1530 | 1570 | ||
@@ -1537,7 +1577,7 @@ def _shutdown(): | |||
| 1537 | 1577 | atexit_call() | |
| 1538 | 1578 | ||
| 1539 | 1579 | if _is_main_interpreter(): | |
| 1540 | - _main_thread._handle._set_done() | ||
| 1580 | + _main_thread._os_thread_handle._set_done() | ||
| 1541 | 1581 | ||
| 1542 | 1582 | # Wait for all non-daemon threads to exit. | |
| 1543 | 1583 | _thread_shutdown() | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1301,6 +1301,8 @@ mod sys { | |||
| 1301 | 1301 | safe_path: bool, | |
| 1302 | 1302 | /// -X warn_default_encoding, PYTHONWARNDEFAULTENCODING | |
| 1303 | 1303 | warn_default_encoding: u8, | |
| 1304 | + /// -X thread_inherit_context, whether new threads inherit context from parent | ||
| 1305 | + thread_inherit_context: bool, | ||
| 1304 | 1306 | } | |
| 1305 | 1307 | ||
| 1306 | 1308 | impl FlagsData { | |
@@ -1324,6 +1326,7 @@ mod sys { | |||
| 1324 | 1326 | int_max_str_digits: settings.int_max_str_digits, | |
| 1325 | 1327 | safe_path: settings.safe_path, | |
| 1326 | 1328 | warn_default_encoding: settings.warn_default_encoding as u8, | |
| 1329 | + thread_inherit_context: settings.thread_inherit_context, | ||
| 1327 | 1330 | } | |
| 1328 | 1331 | } | |
| 1329 | 1332 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -88,6 +88,9 @@ pub struct Settings { | |||
| 88 | 88 | /// -X warn_default_encoding, PYTHONWARNDEFAULTENCODING | |
| 89 | 89 | pub warn_default_encoding: bool, | |
| 90 | 90 | ||
| 91 | + /// -X thread_inherit_context, whether new threads inherit context from parent | ||
| 92 | + pub thread_inherit_context: bool, | ||
| 93 | + | ||
| 91 | 94 | /// -i | |
| 92 | 95 | pub inspect: bool, | |
| 93 | 96 | ||
@@ -190,6 +193,7 @@ impl Default for Settings { | |||
| 190 | 193 | isolated: false, | |
| 191 | 194 | dev_mode: false, | |
| 192 | 195 | warn_default_encoding: false, | |
| 196 | + thread_inherit_context: false, | ||
| 193 | 197 | warnoptions: vec![], | |
| 194 | 198 | path_list: vec![], | |
| 195 | 199 | argv: vec![], | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -285,6 +285,20 @@ pub fn parse_opts() -> Result<(Settings, RunMode), lexopt::Error> { | |||
| 285 | 285 | } | |
| 286 | 286 | }; | |
| 287 | 287 | } | |
| 288 | + "thread_inherit_context" => { | ||
| 289 | + settings.thread_inherit_context = match value { | ||
| 290 | + Some("1") => true, | ||
| 291 | + Some("0") => false, | ||
| 292 | + _ => { | ||
| 293 | + error!( | ||
| 294 | + "Fatal Python error: config_init_thread_inherit_context: \ | ||
| 295 | + -X thread_inherit_context=n: n is missing or invalid\n\ | ||
| 296 | + Python runtime state: preinitialized" | ||
| 297 | + ); | ||
| 298 | + std::process::exit(1); | ||
| 299 | + } | ||
| 300 | + }; | ||
| 301 | + } | ||
| 288 | 302 | _ => {} | |
| 289 | 303 | } | |
| 290 | 304 | (name, value.map(str::to_owned)) | |
@@ -297,6 +311,20 @@ pub fn parse_opts() -> Result<(Settings, RunMode), lexopt::Error> { | |||
| 297 | 311 | if env_bool("PYTHONNODEBUGRANGES") { | |
| 298 | 312 | settings.code_debug_ranges = false; | |
| 299 | 313 | } | |
| 314 | + if let Some(val) = get_env("PYTHON_THREAD_INHERIT_CONTEXT") { | ||
| 315 | + settings.thread_inherit_context = match val.to_str() { | ||
| 316 | + Some("1") => true, | ||
| 317 | + Some("0") => false, | ||
| 318 | + _ => { | ||
| 319 | + error!( | ||
| 320 | + "Fatal Python error: config_init_thread_inherit_context: \ | ||
| 321 | + PYTHON_THREAD_INHERIT_CONTEXT=N: N is missing or invalid\n\ | ||
| 322 | + Python runtime state: preinitialized" | ||
| 323 | + ); | ||
| 324 | + std::process::exit(1); | ||
| 325 | + } | ||
| 326 | + }; | ||
| 327 | + } | ||
| 300 | 328 | ||
| 301 | 329 | // Parse PYTHONIOENCODING=encoding[:errors] | |
| 302 | 330 | if let Some(val) = get_env("PYTHONIOENCODING") | |
| Back | FazBrowse Home | New Git URL |
0 commit comments