Skip to content

Commit 2ae2547

Browse files
committed
Pyo3: Type signatures
1 parent 7f83e4f commit 2ae2547

4 files changed

Lines changed: 30 additions & 20 deletions

File tree

‎libs/opsqueue_python/src/async_util.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@ where
3636
}
3737

3838
/// Version of `future_into_py` that uses the `TokioRuntimeThatIsInScope`
39-
pub fn future_into_py<T, F>(py: Python<'_>, fut: F) -> PyResult<Bound<'_, PyAny>>
39+
pub(crate) fn future_into_py<T, F>(py: Python<'_>, fut: F) -> PyResult<Bound<'_, PyAny>>
4040
where
4141
F: Future<Output = PyResult<T>> + Send + 'static,
4242
T: for<'py> IntoPyObject<'py> + Send + 'static,

‎libs/opsqueue_python/src/common.rs‎

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ pub struct SubmissionId {
3333
#[pymethods]
3434
impl SubmissionId {
3535
#[new]
36+
#[pyo3(signature = (id))]
3637
fn new(id: u64) -> CPyResult<Self, TryFromIntError> {
3738
let _is_inner_valid =
3839
opsqueue::common::submission::SubmissionId::try_from(id).map_err(CError)?;
@@ -69,6 +70,7 @@ pub struct ChunkIndex {
6970
#[pymethods]
7071
impl ChunkIndex {
7172
#[new]
73+
#[pyo3(signature = (id))]
7274
fn new(id: u64) -> CPyResult<Self, TryFromIntError> {
7375
let _is_inner_valid = opsqueue::common::chunk::ChunkIndex::new(id).map_err(CError)?;
7476
Ok(ChunkIndex { id })
@@ -230,7 +232,7 @@ impl Chunk {
230232
/// # Errors
231233
///
232234
/// Returns an error if fetching chunk bytes from object storage fails.
233-
pub async fn from_internal(
235+
pub(crate) async fn from_internal(
234236
c: chunk::Chunk,
235237
s: submission::Submission,
236238
object_store_client: &ObjectStoreClient,
@@ -285,7 +287,7 @@ pub struct ChunkFailed {
285287

286288
impl ChunkFailed {
287289
#[must_use]
288-
pub fn from_internal(c: chunk::ChunkFailed, _s: &submission::SubmissionFailed) -> Self {
290+
pub(crate) fn from_internal(c: chunk::ChunkFailed, _s: &submission::SubmissionFailed) -> Self {
289291
ChunkFailed {
290292
submission_id: c.submission_id.into(),
291293
chunk_index: c.chunk_index.into(),
@@ -561,7 +563,7 @@ impl SubmissionNotCancellable {
561563
///
562564
/// Returns the underlying future error, or a fatal Python exception when
563565
/// an interrupt signal is detected.
564-
pub async fn run_unless_interrupted<T, E>(
566+
pub(crate) async fn run_unless_interrupted<T, E>(
565567
future: impl IntoFuture<Output = Result<T, E>>,
566568
) -> Result<T, E>
567569
where
@@ -573,7 +575,7 @@ where
573575
}
574576
}
575577

576-
pub async fn check_signals_in_background() -> FatalPythonException {
578+
async fn check_signals_in_background() -> FatalPythonException {
577579
loop {
578580
tokio::time::sleep(SIGNAL_CHECK_INTERVAL).await;
579581
let res = Python::attach(|py| {
@@ -615,7 +617,7 @@ pub async fn check_signals_in_background() -> FatalPythonException {
615617
///
616618
/// Panics if creating the Tokio runtime fails.
617619
#[must_use]
618-
pub fn start_runtime() -> Arc<tokio::runtime::Runtime> {
620+
pub(crate) fn start_runtime() -> Arc<tokio::runtime::Runtime> {
619621
let runtime = tokio::runtime::Builder::new_multi_thread()
620622
.worker_threads(1)
621623
.enable_all()
@@ -635,7 +637,7 @@ pub fn start_runtime() -> Arc<tokio::runtime::Runtime> {
635637
/// # Panics
636638
///
637639
/// Panics if formatting a Python traceback fails.
638-
pub fn format_pyerr(err: &PyErr) -> String {
640+
pub(crate) fn format_pyerr(err: &PyErr) -> String {
639641
Python::attach(|py| {
640642
let msg: Option<String> = (|| {
641643
let traceback = err.traceback(py)?;

‎libs/opsqueue_python/src/consumer.rs‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -91,12 +91,13 @@ impl ConsumerClient {
9191
)
9292
}
9393

94-
#[allow(clippy::type_complexity)]
9594
/// Reserve up to `max` chunks from the queue.
9695
///
9796
/// # Errors
9897
///
9998
/// Returns an error if reservation fails or if chunk bytes cannot be retrieved.
99+
#[allow(clippy::type_complexity)]
100+
#[pyo3(signature = (max, strategy))]
100101
pub fn reserve_chunks(
101102
&self,
102103
py: Python<'_>,
@@ -114,12 +115,12 @@ impl ConsumerClient {
114115
py.detach(|| self.reserve_chunks_gilless(max, strategy.into()))
115116
}
116117

117-
#[pyo3(signature = (submission_id, submission_prefix, chunk_index, output_content))]
118118
/// Complete a chunk and optionally upload output content to object storage.
119119
///
120120
/// # Errors
121121
///
122122
/// Returns an error if upload or completion request fails.
123+
#[pyo3(signature = (submission_id, submission_prefix, chunk_index, output_content))]
123124
pub fn complete_chunk(
124125
&self,
125126
py: Python<'_>,
@@ -145,12 +146,12 @@ impl ConsumerClient {
145146
})
146147
}
147148

148-
#[pyo3(signature = (submission_id, submission_prefix, chunk_index, failure))]
149149
/// Mark a chunk as failed.
150150
///
151151
/// # Errors
152152
///
153153
/// Returns an error if the failure report cannot be submitted.
154+
#[pyo3(signature = (submission_id, submission_prefix, chunk_index, failure))]
154155
pub fn fail_chunk(
155156
&self,
156157
py: Python<'_>,
@@ -165,6 +166,7 @@ impl ConsumerClient {
165166
}
166167

167168
#[allow(clippy::type_complexity)]
169+
#[pyo3(signature = (strategy, fun))]
168170
pub fn run_per_chunk(
169171
&self,
170172
strategy: &Strategy,
@@ -316,7 +318,7 @@ impl ConsumerClient {
316318
/// # Errors
317319
///
318320
/// Returns an error if the failure report cannot be submitted.
319-
pub fn fail_chunk_gilless(
321+
fn fail_chunk_gilless(
320322
&self,
321323
submission_id: SubmissionId,
322324
_submission_prefix: Option<String>,

‎libs/opsqueue_python/src/producer.rs‎

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -135,6 +135,7 @@ impl ProducerClient {
135135
///
136136
/// Returns an error if the submission cannot be cancelled or if the request fails.
137137
#[allow(clippy::result_large_err, clippy::type_complexity)]
138+
#[pyo3(signature = (id))]
138139
pub fn cancel_submission(
139140
&self,
140141
py: Python<'_>,
@@ -168,6 +169,7 @@ impl ProducerClient {
168169
/// # Errors
169170
///
170171
/// Returns an error if contacting the server fails.
172+
#[pyo3(signature = (id))]
171173
pub fn get_submission_status(
172174
&self,
173175
py: Python<'_>,
@@ -193,6 +195,7 @@ impl ProducerClient {
193195
/// # Errors
194196
///
195197
/// Returns an error if contacting the server fails.
198+
#[pyo3(signature = (prefix))]
196199
pub fn lookup_submission_id_by_prefix(
197200
&self,
198201
py: Python<'_>,
@@ -216,6 +219,7 @@ impl ProducerClient {
216219
///
217220
/// Returns an error if too many submissions match or if the request fails.
218221
#[allow(clippy::needless_pass_by_value)]
222+
#[pyo3(signature = (strategic_metadata))]
219223
pub fn lookup_submission_ids_by_strategic_metadata(
220224
&self,
221225
py: Python<'_>,
@@ -246,25 +250,24 @@ impl ProducerClient {
246250
/// # Errors
247251
///
248252
/// Returns an error if submission insertion fails.
249-
#[pyo3(signature = (chunk_contents, metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))]
253+
#[pyo3(signature = (chunk_contents, metadata=None, strategic_metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))]
250254
pub fn insert_submission_direct(
251255
&self,
252256
py: Python<'_>,
253257
chunk_contents: Vec<chunk::Content>,
254258
metadata: Option<submission::Metadata>,
259+
strategic_metadata: Option<StrategicMetadataMap>,
255260
chunk_size: Option<u64>,
256261
otel_trace_carrier: CarrierMap,
257262
) -> CPyResult<SubmissionId, E<FatalPythonException, InternalProducerClientError>> {
258-
let strategic_metadata = std::collections::HashMap::default();
259-
260263
py.detach(|| {
261264
let submission = opsqueue::producer::InsertSubmission {
262265
chunk_size: chunk_size.map(|n| chunk::ChunkSize(n.cast_signed())),
263266
chunk_contents: ChunkContents::Direct {
264267
contents: chunk_contents,
265268
},
266269
metadata,
267-
strategic_metadata,
270+
strategic_metadata: strategic_metadata.unwrap_or_default(),
268271
};
269272
self.block_unless_interrupted(async move {
270273
self.client
@@ -276,13 +279,13 @@ impl ProducerClient {
276279
})
277280
}
278281

279-
#[pyo3(signature = (chunk_contents, metadata=None, strategic_metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))]
280-
#[allow(clippy::type_complexity)]
281282
/// Insert submission chunks via object storage and enqueue the submission.
282283
///
283284
/// # Errors
284285
///
285286
/// Returns an error if chunk upload or submission insertion fails.
287+
#[allow(clippy::type_complexity)]
288+
#[pyo3(signature = (chunk_contents, metadata=None, strategic_metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))]
286289
pub fn insert_submission_chunks(
287290
&self,
288291
py: Python<'_>,
@@ -343,12 +346,13 @@ impl ProducerClient {
343346
})
344347
}
345348

346-
#[allow(clippy::result_large_err, clippy::type_complexity)]
347349
/// Try streaming completed submission chunks without waiting.
348350
///
349351
/// # Errors
350352
///
351353
/// Returns an error if the submission is incomplete/failed or if the request fails.
354+
#[allow(clippy::result_large_err, clippy::type_complexity)]
355+
#[pyo3(signature = (id))]
352356
pub fn try_stream_completed_submission_chunks(
353357
&self,
354358
py: Python<'_>,
@@ -376,13 +380,13 @@ impl ProducerClient {
376380
})
377381
}
378382

379-
#[pyo3(signature = (chunk_contents, metadata=None, strategic_metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))]
380-
#[allow(clippy::result_large_err, clippy::type_complexity)]
381383
/// Submit chunks and then stream the completed output chunks.
382384
///
383385
/// # Errors
384386
///
385387
/// Returns an error if upload, submission creation, or streaming fails.
388+
#[allow(clippy::result_large_err, clippy::type_complexity)]
389+
#[pyo3(signature = (chunk_contents, metadata=None, strategic_metadata=None, chunk_size=None, otel_trace_carrier=CarrierMap::default()))]
386390
pub fn run_submission_chunks(
387391
&self,
388392
py: Python<'_>,
@@ -438,6 +442,7 @@ impl ProducerClient {
438442
///
439443
/// Returns an error if polling or output streaming fails.
440444
#[allow(clippy::result_large_err, clippy::type_complexity)]
445+
#[pyo3(signature = (submission_id))]
441446
pub fn blocking_stream_completed_submission_chunks(
442447
&self,
443448
py: Python<'_>,
@@ -462,6 +467,7 @@ impl ProducerClient {
462467
/// # Errors
463468
///
464469
/// Returns a Python error if creating the awaitable fails.
470+
#[pyo3(signature = (submission_id))]
465471
pub fn async_stream_completed_submission_chunks<'p>(
466472
&self,
467473
py: Python<'p>,

0 commit comments

Comments
 (0)