Python 3 compatibility. (#121)

* Make common module Python 3 compatible.

* Make plasma module Python 3 compatible.

* Make photon module Python 3 compatible.

* Make numbuf module Python 3 compatible.

* Remaining changes for Python 3 compatibility.

* Test Python 3 in Travis.

* Fixes.
This commit is contained in:
Robert Nishihara
2016-12-16 14:40:37 -08:00
parent 946242929f
commit 79dd1815a2
42 changed files with 652 additions and 302 deletions
+76 -62
View File
@@ -1,4 +1,5 @@
#include <Python.h>
#include "bytesobject.h"
#include "node.h"
#include "common.h"
@@ -17,12 +18,19 @@ PyObject *pickle_dumps = NULL;
PyObject *pickle_protocol = NULL;
void init_pickle_module(void) {
/* For Python 3 this needs to be "_pickle" instead of "cPickle". */
#if PY_MAJOR_VERSION >= 3
pickle_module = PyImport_ImportModule("pickle");
#else
pickle_module = PyImport_ImportModuleNoBlock("cPickle");
pickle_loads = PyString_FromString("loads");
pickle_dumps = PyString_FromString("dumps");
pickle_protocol = PyObject_GetAttrString(pickle_module, "HIGHEST_PROTOCOL");
#endif
CHECK(pickle_module != NULL);
CHECK(PyObject_HasAttrString(pickle_module, "loads"));
CHECK(PyObject_HasAttrString(pickle_module, "dumps"));
CHECK(PyObject_HasAttrString(pickle_module, "HIGHEST_PROTOCOL"));
pickle_loads = PyUnicode_FromString("loads");
pickle_dumps = PyUnicode_FromString("dumps");
pickle_protocol = PyObject_GetAttrString(pickle_module, "HIGHEST_PROTOCOL");
CHECK(pickle_protocol != NULL);
}
/* Define the PyObjectID class. */
@@ -62,8 +70,8 @@ PyObject *PyObjectID_make(object_id object_id) {
static PyObject *PyObjectID_id(PyObject *self) {
PyObjectID *s = (PyObjectID *) self;
return PyString_FromStringAndSize((char *) &s->object_id.id[0],
sizeof(s->object_id.id));
return PyBytes_FromStringAndSize((char *) &s->object_id.id[0],
sizeof(s->object_id.id));
}
static PyObject *PyObjectID_richcompare(PyObjectID *self,
@@ -106,7 +114,7 @@ static PyObject *PyObjectID_richcompare(PyObjectID *self,
static long PyObjectID_hash(PyObjectID *self) {
PyObject *tuple = PyTuple_New(UNIQUE_ID_SIZE);
for (int i = 0; i < UNIQUE_ID_SIZE; ++i) {
PyTuple_SetItem(tuple, i, PyInt_FromLong(self->object_id.id[i]));
PyTuple_SetItem(tuple, i, PyLong_FromLong(self->object_id.id[i]));
}
long hash = PyObject_Hash(tuple);
Py_XDECREF(tuple);
@@ -120,7 +128,7 @@ static PyObject *PyObjectID_repr(PyObjectID *self) {
UT_string *repr;
utstring_new(repr);
utstring_printf(repr, "ObjectID(%s)", hex_id);
PyObject *result = PyString_FromString(utstring_body(repr));
PyObject *result = PyUnicode_FromString(utstring_body(repr));
utstring_free(repr);
return result;
}
@@ -144,7 +152,7 @@ static PyMemberDef PyObjectID_members[] = {
};
PyTypeObject PyObjectIDType = {
PyObject_HEAD_INIT(NULL) 0, /* ob_size */
PyVarObject_HEAD_INIT(NULL, 0) /* ob_size */
"common.ObjectID", /* tp_name */
sizeof(PyObjectID), /* tp_basicsize */
0, /* tp_itemsize */
@@ -203,17 +211,17 @@ static int PyTask_init(PyTask *self, PyObject *args, PyObject *kwds) {
&parent_task_id, &parent_counter)) {
return -1;
}
size_t size = PyList_Size(arguments);
Py_ssize_t size = PyList_Size(arguments);
/* Determine the size of pass by value data in bytes. */
size_t value_data_bytes = 0;
for (size_t i = 0; i < size; ++i) {
Py_ssize_t value_data_bytes = 0;
for (Py_ssize_t i = 0; i < size; ++i) {
PyObject *arg = PyList_GetItem(arguments, i);
if (!PyObject_IsInstance(arg, (PyObject *) &PyObjectIDType)) {
CHECK(pickle_module != NULL);
CHECK(pickle_dumps != NULL);
PyObject *data = PyObject_CallMethodObjArgs(pickle_module, pickle_dumps,
arg, pickle_protocol, NULL);
value_data_bytes += PyString_Size(data);
value_data_bytes += PyBytes_Size(data);
utarray_push_back(val_repr_ptrs, &data);
}
}
@@ -223,15 +231,17 @@ static int PyTask_init(PyTask *self, PyObject *args, PyObject *kwds) {
start_construct_task_spec(parent_task_id, parent_counter, function_id,
size, num_returns, value_data_bytes);
/* Add the task arguments. */
for (size_t i = 0; i < size; ++i) {
for (Py_ssize_t i = 0; i < size; ++i) {
PyObject *arg = PyList_GetItem(arguments, i);
if (PyObject_IsInstance(arg, (PyObject *) &PyObjectIDType)) {
task_args_add_ref(self->spec, ((PyObjectID *) arg)->object_id);
} else {
PyObject *data =
*((PyObject **) utarray_eltptr(val_repr_ptrs, val_repr_index));
task_args_add_val(self->spec, (uint8_t *) PyString_AS_STRING(data),
PyString_GET_SIZE(data));
/* We do this check because we cast a signed int to an unsigned int. */
CHECK(val_repr_index >= 0);
PyObject *data = *((PyObject **) utarray_eltptr(
val_repr_ptrs, (uint64_t) val_repr_index));
task_args_add_val(self->spec, (uint8_t *) PyBytes_AS_STRING(data),
PyBytes_GET_SIZE(data));
Py_DECREF(data);
val_repr_index += 1;
}
@@ -269,8 +279,8 @@ static PyObject *PyTask_arguments(PyObject *self) {
CHECK(pickle_module != NULL);
CHECK(pickle_loads != NULL);
PyObject *str =
PyString_FromStringAndSize((char *) task_arg_val(task, i),
(Py_ssize_t) task_arg_length(task, i));
PyBytes_FromStringAndSize((char *) task_arg_val(task, i),
(Py_ssize_t) task_arg_length(task, i));
PyObject *val =
PyObject_CallMethodObjArgs(pickle_module, pickle_loads, str, NULL);
Py_XDECREF(str);
@@ -304,44 +314,44 @@ static PyMethodDef PyTask_methods[] = {
};
PyTypeObject PyTaskType = {
PyObject_HEAD_INIT(NULL) 0, /* ob_size */
"task.Task", /* tp_name */
sizeof(PyTask), /* tp_basicsize */
0, /* tp_itemsize */
(destructor) PyTask_dealloc, /* tp_dealloc */
0, /* tp_print */
0, /* tp_getattr */
0, /* tp_setattr */
0, /* tp_compare */
0, /* tp_repr */
0, /* tp_as_number */
0, /* tp_as_sequence */
0, /* tp_as_mapping */
0, /* tp_hash */
0, /* tp_call */
0, /* tp_str */
0, /* tp_getattro */
0, /* tp_setattro */
0, /* tp_as_buffer */
Py_TPFLAGS_DEFAULT, /* tp_flags */
"Task object", /* tp_doc */
0, /* tp_traverse */
0, /* tp_clear */
0, /* tp_richcompare */
0, /* tp_weaklistoffset */
0, /* tp_iter */
0, /* tp_iternext */
PyTask_methods, /* tp_methods */
0, /* tp_members */
0, /* tp_getset */
0, /* tp_base */
0, /* tp_dict */
0, /* tp_descr_get */
0, /* tp_descr_set */
0, /* tp_dictoffset */
(initproc) PyTask_init, /* tp_init */
0, /* tp_alloc */
PyType_GenericNew, /* tp_new */
PyVarObject_HEAD_INIT(NULL, 0) /* ob_size */
"task.Task", /* tp_name */
sizeof(PyTask), /* tp_basicsize */
0, /* tp_itemsize */
(destructor) PyTask_dealloc, /* tp_dealloc */
0, /* tp_print */
0, /* tp_getattr */
0, /* tp_setattr */
0, /* tp_compare */
0, /* tp_repr */
0, /* tp_as_number */
0, /* tp_as_sequence */
0, /* tp_as_mapping */
0, /* tp_hash */
0, /* tp_call */
0, /* tp_str */
0, /* tp_getattro */
0, /* tp_setattro */
0, /* tp_as_buffer */
Py_TPFLAGS_DEFAULT, /* tp_flags */
"Task object", /* tp_doc */
0, /* tp_traverse */
0, /* tp_clear */
0, /* tp_richcompare */
0, /* tp_weaklistoffset */
0, /* tp_iter */
0, /* tp_iternext */
PyTask_methods, /* tp_methods */
0, /* tp_members */
0, /* tp_getset */
0, /* tp_base */
0, /* tp_dict */
0, /* tp_descr_get */
0, /* tp_descr_set */
0, /* tp_dictoffset */
(initproc) PyTask_init, /* tp_init */
0, /* tp_alloc */
PyType_GenericNew, /* tp_new */
};
/* Create a PyTask from a C struct. The resulting PyTask takes ownership of the
@@ -358,6 +368,10 @@ PyObject *PyTask_make(task_spec *task_spec) {
#define SIZE_LIMIT 100
#define NUM_ELEMENTS_LIMIT 1000
#if PY_MAJOR_VERSION >= 3
#define PyInt_Check PyLong_Check
#endif
/**
* This method checks if a Python object is sufficiently simple that it can be
* serialized and passed by value as an argument to a task (without being put in
@@ -382,8 +396,8 @@ int is_simple_value(PyObject *value, int *num_elements_contained) {
value == Py_True || PyFloat_Check(value) || value == Py_None) {
return 1;
}
if (PyString_CheckExact(value)) {
*num_elements_contained += PyString_Size(value);
if (PyBytes_CheckExact(value)) {
*num_elements_contained += PyBytes_Size(value);
return (*num_elements_contained < NUM_ELEMENTS_LIMIT);
}
if (PyUnicode_CheckExact(value)) {
@@ -391,7 +405,7 @@ int is_simple_value(PyObject *value, int *num_elements_contained) {
return (*num_elements_contained < NUM_ELEMENTS_LIMIT);
}
if (PyList_CheckExact(value) && PyList_Size(value) < SIZE_LIMIT) {
for (size_t i = 0; i < PyList_Size(value); ++i) {
for (Py_ssize_t i = 0; i < PyList_Size(value); ++i) {
if (!is_simple_value(PyList_GetItem(value, i), num_elements_contained)) {
return 0;
}
@@ -410,7 +424,7 @@ int is_simple_value(PyObject *value, int *num_elements_contained) {
return (*num_elements_contained < NUM_ELEMENTS_LIMIT);
}
if (PyTuple_CheckExact(value) && PyTuple_Size(value) < SIZE_LIMIT) {
for (size_t i = 0; i < PyTuple_Size(value); ++i) {
for (Py_ssize_t i = 0; i < PyTuple_Size(value); ++i) {
if (!is_simple_value(PyTuple_GetItem(value, i), num_elements_contained)) {
return 0;
}
+41 -5
View File
@@ -10,21 +10,53 @@ static PyMethodDef common_methods[] = {
{NULL} /* Sentinel */
};
#if PY_MAJOR_VERSION >= 3
static struct PyModuleDef moduledef = {
PyModuleDef_HEAD_INIT,
"common", /* m_name */
"A module for common types. This is used for testing.", /* m_doc */
0, /* m_size */
common_methods, /* m_methods */
NULL, /* m_reload */
NULL, /* m_traverse */
NULL, /* m_clear */
NULL, /* m_free */
};
#endif
#if PY_MAJOR_VERSION >= 3
#define INITERROR return NULL
#else
#define INITERROR return
#endif
#ifndef PyMODINIT_FUNC /* declarations for DLL import/export */
#define PyMODINIT_FUNC void
#endif
PyMODINIT_FUNC initcommon(void) {
#if PY_MAJOR_VERSION >= 3
#define MOD_INIT(name) PyMODINIT_FUNC PyInit_##name(void)
#else
#define MOD_INIT(name) PyMODINIT_FUNC init##name(void)
#endif
MOD_INIT(common) {
PyObject *m;
if (PyType_Ready(&PyTaskType) < 0)
return;
if (PyType_Ready(&PyTaskType) < 0) {
INITERROR;
}
if (PyType_Ready(&PyObjectIDType) < 0)
return;
if (PyType_Ready(&PyObjectIDType) < 0) {
INITERROR;
}
#if PY_MAJOR_VERSION >= 3
m = PyModule_Create(&moduledef);
#else
m = Py_InitModule3("common", common_methods,
"A module for common types. This is used for testing.");
#endif
init_pickle_module();
@@ -38,4 +70,8 @@ PyMODINIT_FUNC initcommon(void) {
CommonError = PyErr_NewException(common_error, NULL, NULL);
Py_INCREF(CommonError);
PyModule_AddObject(m, "common_error", CommonError);
#if PY_MAJOR_VERSION >= 3
return m;
#endif
}
+5 -5
View File
@@ -63,11 +63,11 @@ class TestGlobalStateStore(unittest.TestCase):
self.redis.execute_command("RAY.OBJECT_TABLE_ADD", "object_id1", 1, "hash1", "manager_id1")
self.redis.execute_command("RAY.OBJECT_TABLE_ADD", "object_id1", 1, "hash1", "manager_id2")
response = self.redis.execute_command("RAY.OBJECT_TABLE_LOOKUP", "object_id1")
self.assertEqual(set(response), {"manager_id1", "manager_id2"})
self.assertEqual(set(response), {b"manager_id1", b"manager_id2"})
# Add a manager that already exists again and try again.
self.redis.execute_command("RAY.OBJECT_TABLE_ADD", "object_id1", 1, "hash1", "manager_id2")
response = self.redis.execute_command("RAY.OBJECT_TABLE_LOOKUP", "object_id1")
self.assertEqual(set(response), {"manager_id1", "manager_id2"})
self.assertEqual(set(response), {b"manager_id1", b"manager_id2"})
def testObjectTableSubscribe(self):
p = self.redis.pubsub()
@@ -77,7 +77,7 @@ class TestGlobalStateStore(unittest.TestCase):
# Receive the acknowledgement message.
self.assertEqual(p.get_message()["data"], 1)
# Receive the actual data.
self.assertEqual(p.get_message()["data"], "MANAGERS manager_id1")
self.assertEqual(p.get_message()["data"], b"MANAGERS manager_id1")
def testResultTableAddAndLookup(self):
response = self.redis.execute_command("RAY.RESULT_TABLE_LOOKUP", "object_id1")
@@ -87,10 +87,10 @@ class TestGlobalStateStore(unittest.TestCase):
self.assertEqual(set(response), set([]))
self.redis.execute_command("RAY.RESULT_TABLE_ADD", "object_id1", "task_id1")
response = self.redis.execute_command("RAY.RESULT_TABLE_LOOKUP", "object_id1")
self.assertEqual(response, "task_id1")
self.assertEqual(response, b"task_id1")
self.redis.execute_command("RAY.RESULT_TABLE_ADD", "object_id2", "task_id2")
response = self.redis.execute_command("RAY.RESULT_TABLE_LOOKUP", "object_id2")
self.assertEqual(response, "task_id2")
self.assertEqual(response, b"task_id2")
if __name__ == "__main__":
unittest.main(verbosity=2)
+13 -8
View File
@@ -4,6 +4,7 @@ from __future__ import print_function
import numpy as np
import pickle
import sys
import unittest
import common
@@ -20,10 +21,13 @@ def random_task_id():
return common.ObjectID(np.random.bytes(ID_SIZE))
BASE_SIMPLE_OBJECTS = [
0, 1, 100000, 0L, 1L, 100000L, 1L << 100, 0.0, 0.5, 0.9, 100000.1, (), [], {},
0, 1, 100000, 0.0, 0.5, 0.9, 100000.1, (), [], {},
"", 990 * "h", u"", 990 * u"h"
]
if sys.version_info < (3, 0):
BASE_SIMPLE_OBJECTS += [long(0), long(1), long(100000), long(1 << 100)]
LIST_SIMPLE_OBJECTS = [[obj] for obj in BASE_SIMPLE_OBJECTS]
TUPLE_SIMPLE_OBJECTS = [(obj,) for obj in BASE_SIMPLE_OBJECTS]
DICT_SIMPLE_OBJECTS = [{(): obj} for obj in BASE_SIMPLE_OBJECTS]
@@ -85,16 +89,17 @@ class TestObjectID(unittest.TestCase):
self.assertRaises(Exception, lambda : pickling.dumps(h))
def test_equality_comparisons(self):
x1 = common.ObjectID(ID_SIZE * "a")
x2 = common.ObjectID(ID_SIZE * "a")
y1 = common.ObjectID(ID_SIZE * "b")
y2 = common.ObjectID(ID_SIZE * "b")
x1 = common.ObjectID(ID_SIZE * b"a")
x2 = common.ObjectID(ID_SIZE * b"a")
y1 = common.ObjectID(ID_SIZE * b"b")
y2 = common.ObjectID(ID_SIZE * b"b")
self.assertEqual(x1, x2)
self.assertEqual(y1, y2)
self.assertNotEqual(x1, y1)
object_ids1 = [common.ObjectID(ID_SIZE * chr(i)) for i in range(256)]
object_ids2 = [common.ObjectID(ID_SIZE * chr(i)) for i in range(256)]
random_strings = [np.random.bytes(ID_SIZE) for _ in range(256)]
object_ids1 = [common.ObjectID(random_strings[i]) for i in range(256)]
object_ids2 = [common.ObjectID(random_strings[i]) for i in range(256)]
self.assertEqual(len(set(object_ids1)), 256)
self.assertEqual(len(set(object_ids1 + object_ids2)), 256)
self.assertEqual(set(object_ids1), set(object_ids2))
@@ -122,7 +127,7 @@ class TestTask(unittest.TestCase):
10 * ["a"],
100 * ["a"],
1000 * ["a"],
[1, 1.3, 2L, 1L << 100, "hi", u"hi", [1, 2]],
[1, 1.3, 2, 1 << 100, "hi", u"hi", [1, 2]],
object_ids[:1],
object_ids[:2],
object_ids[:3],
+2 -2
View File
@@ -102,7 +102,7 @@ class TestGlobalScheduler(unittest.TestCase):
self.assertLessEqual(len(task_entries), 1)
if len(task_entries) == 1:
task_contents = self.redis_client.hgetall(task_entries[0])
task_status = int(task_contents["state"])
task_status = int(task_contents[b"state"])
self.assertTrue(task_status in [TASK_STATUS_WAITING, TASK_STATUS_SCHEDULED])
if task_status == TASK_STATUS_SCHEDULED:
break
@@ -120,7 +120,7 @@ class TestGlobalScheduler(unittest.TestCase):
self.assertLessEqual(len(task_entries), num_tasks + 1)
if len(task_entries) == num_tasks + 1:
task_contents = [self.redis_client.hgetall(task_entries[i]) for i in range(len(task_entries))]
task_statuses = [int(contents["state"]) for contents in task_contents]
task_statuses = [int(contents[b"state"]) for contents in task_contents]
self.assertTrue(all([status in [TASK_STATUS_WAITING, TASK_STATUS_SCHEDULED] for status in task_statuses]))
if all([status == TASK_STATUS_SCHEDULED for status in task_statuses]):
break
+4 -2
View File
@@ -19,8 +19,9 @@ message(STATUS "PYTHON_INCLUDE_DIRS: " ${PYTHON_INCLUDE_DIRS})
execute_process(COMMAND ${CUSTOM_PYTHON_EXECUTABLE} -c "import sys; print(sys.exec_prefix)"
OUTPUT_VARIABLE PYTHON_PREFIX OUTPUT_STRIP_TRAILING_WHITESPACE)
message(STATUS "PYTHON_PREFIX: " ${PYTHON_PREFIX})
# The name ending in "m" is for miniconda.
FIND_LIBRARY(PYTHON_LIBRARIES
NAMES ${PYTHON_LIBRARY_NAME}
NAMES "${PYTHON_LIBRARY_NAME}" "${PYTHON_LIBRARY_NAME}m"
HINTS "${PYTHON_PREFIX}"
PATH_SUFFIXES "lib" "libs"
NO_DEFAULT_PATH)
@@ -29,8 +30,9 @@ message(STATUS "PYTHON_LIBRARIES: " ${PYTHON_LIBRARIES})
# the Python include directories.
if(NOT PYTHON_LIBRARIES)
message(STATUS "Failed to find PYTHON_LIBRARIES near the Python executable, so now looking near the Python include directories.")
# The name ending in "m" is for miniconda.
FIND_LIBRARY(PYTHON_LIBRARIES
NAMES ${PYTHON_LIBRARY_NAME}
NAMES "${PYTHON_LIBRARY_NAME}" "${PYTHON_LIBRARY_NAME}m"
HINTS "${PYTHON_INCLUDE_DIRS}/../.."
PATH_SUFFIXES "lib" "libs"
NO_DEFAULT_PATH)
+2 -2
View File
@@ -2,5 +2,5 @@ from __future__ import absolute_import
from __future__ import division
from __future__ import print_function
from libphoton import *
from photon_services import *
from .libphoton import *
from .photon_services import *
+46 -11
View File
@@ -73,7 +73,7 @@ static PyMethodDef PyPhotonClient_methods[] = {
};
static PyTypeObject PyPhotonClientType = {
PyObject_HEAD_INIT(NULL) 0, /* ob_size */
PyVarObject_HEAD_INIT(NULL, 0) /* ob_size */
"photon.PhotonClient", /* tp_name */
sizeof(PyPhotonClient), /* tp_basicsize */
0, /* tp_itemsize */
@@ -121,24 +121,55 @@ static PyMethodDef photon_methods[] = {
{NULL} /* Sentinel */
};
#if PY_MAJOR_VERSION >= 3
static struct PyModuleDef moduledef = {
PyModuleDef_HEAD_INIT,
"libphoton", /* m_name */
"A module for the local scheduler.", /* m_doc */
0, /* m_size */
photon_methods, /* m_methods */
NULL, /* m_reload */
NULL, /* m_traverse */
NULL, /* m_clear */
NULL, /* m_free */
};
#endif
#if PY_MAJOR_VERSION >= 3
#define INITERROR return NULL
#else
#define INITERROR return
#endif
#ifndef PyMODINIT_FUNC /* declarations for DLL import/export */
#define PyMODINIT_FUNC void
#endif
PyMODINIT_FUNC initlibphoton(void) {
PyObject *m;
#if PY_MAJOR_VERSION >= 3
#define MOD_INIT(name) PyMODINIT_FUNC PyInit_##name(void)
#else
#define MOD_INIT(name) PyMODINIT_FUNC init##name(void)
#endif
if (PyType_Ready(&PyTaskType) < 0)
return;
MOD_INIT(libphoton) {
if (PyType_Ready(&PyTaskType) < 0) {
INITERROR;
}
if (PyType_Ready(&PyObjectIDType) < 0)
return;
if (PyType_Ready(&PyObjectIDType) < 0) {
INITERROR;
}
if (PyType_Ready(&PyPhotonClientType) < 0)
return;
if (PyType_Ready(&PyPhotonClientType) < 0) {
INITERROR;
}
m = Py_InitModule3("libphoton", photon_methods,
"A module for the local scheduler.");
#if PY_MAJOR_VERSION >= 3
PyObject *m = PyModule_Create(&moduledef);
#else
PyObject *m = Py_InitModule3("libphoton", photon_methods,
"A module for the local scheduler.");
#endif
init_pickle_module();
@@ -155,4 +186,8 @@ PyMODINIT_FUNC initlibphoton(void) {
PhotonError = PyErr_NewException(photon_error, NULL, NULL);
Py_INCREF(PhotonError);
PyModule_AddObject(m, "photon_error", PhotonError);
#if PY_MAJOR_VERSION >= 3
return m;
#endif
}
+1 -1
View File
@@ -74,7 +74,7 @@ class TestPhotonClient(unittest.TestCase):
10 * ["a"],
100 * ["a"],
1000 * ["a"],
[1, 1.3, 2L, 1L << 100, "hi", u"hi", [1, 2]],
[1, 1.3, 1 << 100, "hi", u"hi", [1, 2]],
object_ids[:1],
object_ids[:2],
object_ids[:3],
+5 -3
View File
@@ -19,8 +19,9 @@ message(STATUS "PYTHON_INCLUDE_DIRS: " ${PYTHON_INCLUDE_DIRS})
execute_process(COMMAND ${CUSTOM_PYTHON_EXECUTABLE} -c "import sys; print(sys.exec_prefix)"
OUTPUT_VARIABLE PYTHON_PREFIX OUTPUT_STRIP_TRAILING_WHITESPACE)
message(STATUS "PYTHON_PREFIX: " ${PYTHON_PREFIX})
# The name ending in "m" is for miniconda.
FIND_LIBRARY(PYTHON_LIBRARIES
NAMES ${PYTHON_LIBRARY_NAME}
NAMES "${PYTHON_LIBRARY_NAME}" "${PYTHON_LIBRARY_NAME}m"
HINTS "${PYTHON_PREFIX}"
PATH_SUFFIXES "lib" "libs"
NO_DEFAULT_PATH)
@@ -29,8 +30,9 @@ message(STATUS "PYTHON_LIBRARIES: " ${PYTHON_LIBRARIES})
# the Python include directories.
if(NOT PYTHON_LIBRARIES)
message(STATUS "Failed to find PYTHON_LIBRARIES near the Python executable, so now looking near the Python include directories.")
# The name ending in "m" is for miniconda.
FIND_LIBRARY(PYTHON_LIBRARIES
NAMES ${PYTHON_LIBRARY_NAME}
NAMES "${PYTHON_LIBRARY_NAME}" "${PYTHON_LIBRARY_NAME}m"
HINTS "${PYTHON_INCLUDE_DIRS}/../.."
PATH_SUFFIXES "lib" "libs"
NO_DEFAULT_PATH)
@@ -89,4 +91,4 @@ endif(APPLE)
target_link_libraries(plasma ${COMMON_LIB} ${PYTHON_LIBRARIES})
install(TARGETS plasma DESTINATION ${CMAKE_SOURCE_DIR}/lib/python)
install(TARGETS plasma DESTINATION ${CMAKE_SOURCE_DIR}/plasma)
+5
View File
@@ -0,0 +1,5 @@
from __future__ import absolute_import
from __future__ import division
from __future__ import print_function
from plasma.plasma import *
@@ -5,6 +5,7 @@ from __future__ import print_function
import os
import random
import subprocess
import sys
import time
from . import libplasma
@@ -42,13 +43,18 @@ class PlasmaBuffer(object):
def __getitem__(self, index):
"""Read from the PlasmaBuffer as if it were just a regular buffer."""
return self.buffer[index]
value = self.buffer[index]
if sys.version_info >= (3, 0) and not isinstance(index, slice):
value = chr(value)
return value
def __setitem__(self, index, value):
"""Write to the PlasmaBuffer as if it were just a regular buffer.
This should fail because the buffer should be read only.
"""
if sys.version_info >= (3, 0) and not isinstance(index, slice):
value = ord(value)
self.buffer[index] = value
def __len__(self):
@@ -103,7 +109,7 @@ class PlasmaClient(object):
Exception: An exception is raised if the object could not be created.
"""
# Turn the metadata into the right type.
metadata = bytearray("") if metadata is None else metadata
metadata = bytearray(b"") if metadata is None else metadata
buff = libplasma.create(self.conn, object_id, size, metadata)
return PlasmaBuffer(buff, object_id, self)
@@ -243,7 +249,7 @@ def start_plasma_store(plasma_store_memory=DEFAULT_PLASMA_STORE_MEMORY, use_valg
"""
if use_valgrind and use_profiler:
raise Exception("Cannot use valgrind and profiler at the same time.")
plasma_store_executable = os.path.join(os.path.abspath(os.path.dirname(__file__)), "../../build/plasma_store")
plasma_store_executable = os.path.join(os.path.abspath(os.path.dirname(__file__)), "plasma_store")
plasma_store_name = "/tmp/plasma_store{}".format(random_name())
command = [plasma_store_executable, "-s", plasma_store_name, "-m", str(plasma_store_memory)]
if use_valgrind:
@@ -273,7 +279,7 @@ def start_plasma_manager(store_name, redis_address, num_retries=20, use_valgrind
Raises:
Exception: An exception is raised if the manager could not be started.
"""
plasma_manager_executable = os.path.join(os.path.abspath(os.path.dirname(__file__)), "../../build/plasma_manager")
plasma_manager_executable = os.path.join(os.path.abspath(os.path.dirname(__file__)), "plasma_manager")
plasma_manager_name = "/tmp/plasma_manager{}".format(random_name())
port = None
process = None
+70 -16
View File
@@ -1,4 +1,5 @@
#include <Python.h>
#include "bytesobject.h"
#include "common.h"
#include "io.h"
@@ -17,8 +18,8 @@ static int PyObjectToPlasmaConnection(PyObject *object,
}
static int PyObjectToUniqueID(PyObject *object, object_id *object_id) {
if (PyString_Check(object)) {
memcpy(&object_id->id[0], PyString_AsString(object), UNIQUE_ID_SIZE);
if (PyBytes_Check(object)) {
memcpy(&object_id->id[0], PyBytes_AsString(object), UNIQUE_ID_SIZE);
return 1;
} else {
PyErr_SetString(PyExc_TypeError, "must be a 20 character string");
@@ -75,7 +76,12 @@ PyObject *PyPlasma_create(PyObject *self, PyObject *args) {
"an object with this ID could not be created");
return NULL;
}
#if PY_MAJOR_VERSION >= 3
return PyMemoryView_FromMemory((void *) data, (Py_ssize_t) size, PyBUF_WRITE);
#else
return PyBuffer_FromReadWriteMemory((void *) data, (Py_ssize_t) size);
#endif
}
PyObject *PyPlasma_hash(PyObject *self, PyObject *args) {
@@ -88,7 +94,8 @@ PyObject *PyPlasma_hash(PyObject *self, PyObject *args) {
unsigned char digest[DIGEST_SIZE];
bool success = plasma_compute_object_hash(conn, object_id, digest);
if (success) {
PyObject *digest_string = PyString_FromStringAndSize(digest, DIGEST_SIZE);
PyObject *digest_string =
PyBytes_FromStringAndSize((char *) digest, DIGEST_SIZE);
return digest_string;
} else {
Py_RETURN_NONE;
@@ -132,9 +139,22 @@ PyObject *PyPlasma_get(PyObject *self, PyObject *args) {
plasma_get(conn, object_id, &size, &data, &metadata_size, &metadata);
Py_END_ALLOW_THREADS;
PyObject *t = PyTuple_New(2);
#if PY_MAJOR_VERSION >= 3
PyTuple_SetItem(t, 0, PyMemoryView_FromMemory((void *) data,
(Py_ssize_t) size, PyBUF_READ));
#else
PyTuple_SetItem(t, 0, PyBuffer_FromMemory((void *) data, (Py_ssize_t) size));
PyTuple_SetItem(t, 1, PyByteArray_FromStringAndSize(
(void *) metadata, (Py_ssize_t) metadata_size));
#endif
#if PY_MAJOR_VERSION >= 3
PyTuple_SetItem(
t, 1, PyMemoryView_FromMemory((void *) metadata,
(Py_ssize_t) metadata_size, PyBUF_READ));
#else
PyTuple_SetItem(
t, 1, PyBuffer_FromMemory((void *) metadata, (Py_ssize_t) metadata_size));
#endif
return t;
}
@@ -233,8 +253,8 @@ PyObject *PyPlasma_wait(PyObject *self, PyObject *args) {
if (object_requests[i].status == PLASMA_OBJECT_LOCAL ||
object_requests[i].status == PLASMA_OBJECT_REMOTE) {
PyObject *ready =
PyString_FromStringAndSize((char *) object_requests[i].object_id.id,
sizeof(object_requests[i].object_id));
PyBytes_FromStringAndSize((char *) object_requests[i].object_id.id,
sizeof(object_requests[i].object_id));
PyList_SetItem(ready_ids, num_returned, ready);
PySet_Discard(waiting_ids, ready);
num_returned += 1;
@@ -258,7 +278,7 @@ PyObject *PyPlasma_evict(PyObject *self, PyObject *args) {
return NULL;
}
int64_t evicted_bytes = plasma_evict(conn, (int64_t) num_bytes);
return PyInt_FromLong((long) evicted_bytes);
return PyLong_FromLong((long) evicted_bytes);
}
PyObject *PyPlasma_delete(PyObject *self, PyObject *args) {
@@ -298,7 +318,7 @@ PyObject *PyPlasma_subscribe(PyObject *self, PyObject *args) {
}
int sock = plasma_subscribe(conn);
return PyInt_FromLong(sock);
return PyLong_FromLong(sock);
}
PyObject *PyPlasma_receive_notification(PyObject *self, PyObject *args) {
@@ -320,10 +340,10 @@ PyObject *PyPlasma_receive_notification(PyObject *self, PyObject *args) {
}
/* Construct a tuple from object_info and return. */
PyObject *t = PyTuple_New(3);
PyTuple_SetItem(t, 0, PyString_FromStringAndSize(
PyTuple_SetItem(t, 0, PyBytes_FromStringAndSize(
(char *) object_info.obj_id.id, UNIQUE_ID_SIZE));
PyTuple_SetItem(t, 1, PyInt_FromLong(object_info.data_size));
PyTuple_SetItem(t, 2, PyInt_FromLong(object_info.metadata_size));
PyTuple_SetItem(t, 1, PyLong_FromLong(object_info.data_size));
PyTuple_SetItem(t, 2, PyLong_FromLong(object_info.metadata_size));
return t;
}
@@ -346,7 +366,7 @@ static PyMethodDef plasma_methods[] = {
{"evict", PyPlasma_evict, METH_VARARGS,
"Evict some objects until we recover some number of bytes."},
{"release", PyPlasma_release, METH_VARARGS, "Release the plasma object."},
{"delete", PyPlasma_delete, METH_VARARGS, "Deleta a plasma object."},
{"delete", PyPlasma_delete, METH_VARARGS, "Delete a plasma object."},
{"transfer", PyPlasma_transfer, METH_VARARGS,
"Transfer object to another plasma manager."},
{"subscribe", PyPlasma_subscribe, METH_VARARGS,
@@ -356,11 +376,45 @@ static PyMethodDef plasma_methods[] = {
{NULL} /* Sentinel */
};
#if PY_MAJOR_VERSION >= 3
static struct PyModuleDef moduledef = {
PyModuleDef_HEAD_INIT,
"libplasma", /* m_name */
"A Python client library for plasma.", /* m_doc */
0, /* m_size */
plasma_methods, /* m_methods */
NULL, /* m_reload */
NULL, /* m_traverse */
NULL, /* m_clear */
NULL, /* m_free */
};
#endif
#if PY_MAJOR_VERSION >= 3
#define INITERROR return NULL
#else
#define INITERROR return
#endif
#ifndef PyMODINIT_FUNC /* declarations for DLL import/export */
#define PyMODINIT_FUNC void
#endif
PyMODINIT_FUNC initlibplasma(void) {
Py_InitModule3("libplasma", plasma_methods,
"A Python client library for plasma");
#if PY_MAJOR_VERSION >= 3
#define MOD_INIT(name) PyMODINIT_FUNC PyInit_##name(void)
#else
#define MOD_INIT(name) PyMODINIT_FUNC init##name(void)
#endif
MOD_INIT(libplasma) {
#if PY_MAJOR_VERSION >= 3
PyObject *m = PyModule_Create(&moduledef);
#else
PyObject *m = Py_InitModule3("libplasma", plasma_methods,
"A Python client library for plasma.");
#endif
#if PY_MAJOR_VERSION >= 3
return m;
#endif
}
@@ -9,9 +9,11 @@ import subprocess
class install(_install.install):
def run(self):
subprocess.check_call(["make"], cwd="../../")
subprocess.check_call(["cmake", ".."], cwd="../../build")
subprocess.check_call(["make", "install"], cwd="../../build")
subprocess.check_call(["make"])
subprocess.check_call(["cp", "build/plasma_store", "plasma/plasma_store"])
subprocess.check_call(["cp", "build/plasma_manager", "plasma/plasma_manager"])
subprocess.check_call(["cmake", ".."], cwd="./build")
subprocess.check_call(["make", "install"], cwd="./build")
# Calling _install.install.run(self) does not fetch required packages and
# instead performs an old-style install. See command/install.py in
# setuptools. So, calling do_egg_install() manually here.
@@ -21,7 +23,10 @@ setup(name="Plasma",
version="0.0.1",
description="Plasma client for Python",
packages=find_packages(),
package_data={"plasma": ["libplasma.so"]},
package_data={"plasma": ["plasma_store",
"plasma_manager",
"libplasma.so"],
},
cmdclass={"install": install},
include_package_data=True,
zip_safe=False)
+6 -6
View File
@@ -24,13 +24,13 @@ def random_object_id():
return np.random.bytes(20)
def generate_metadata(length):
metadata = length * ["\x00"]
metadata_buffer = bytearray(length)
if length > 0:
metadata[0] = chr(random.randint(0, 255))
metadata[-1] = chr(random.randint(0, 255))
metadata_buffer[0] = random.randint(0, 255)
metadata_buffer[-1] = random.randint(0, 255)
for _ in range(100):
metadata[random.randint(0, length - 1)] = chr(random.randint(0, 255))
return bytearray("".join(metadata))
metadata_buffer[random.randint(0, length - 1)] = random.randint(0, 255)
return metadata_buffer
def write_to_data_buffer(buff, length):
if length > 0:
@@ -123,7 +123,7 @@ class TestPlasmaClient(unittest.TestCase):
metadata_buffer = self.plasma_client.get_metadata(object_id)
self.assertEqual(len(metadata), len(metadata_buffer))
for i in range(len(metadata)):
self.assertEqual(metadata[i], metadata_buffer[i])
self.assertEqual(chr(metadata[i]), metadata_buffer[i])
def test_create_existing(self):
# This test is partially used to test the code path in which we create an