blob: bbdfc9de9974aeb875a31b4762081a4985c18bcc [file]
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you 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.
use super::Layer;
use super::PythonLayer;
use crate::*;
use opendal::Operator;
/// A layer that limits the number of concurrent operations.
///
/// Notes
/// -----
/// All operators wrapped by this layer will share a common semaphore. This
/// allows you to reuse the same layer across multiple operators, ensuring
/// that the total number of concurrent requests across the entire
/// application does not exceed the limit.
#[pyclass(module = "opendal.layers", extends = Layer, skip_from_py_object)]
#[derive(Clone)]
pub struct ConcurrentLimitLayer(ocore::layers::ConcurrentLimitLayer);
impl PythonLayer for ConcurrentLimitLayer {
fn layer(&self, op: Operator) -> Operator {
op.layer(self.0.clone())
}
}
#[pymethods]
impl ConcurrentLimitLayer {
/// Create a new ConcurrentLimitLayer.
///
/// Parameters
/// ----------
/// limit : int
/// Maximum number of concurrent operations allowed.
///
/// Returns
/// -------
/// ConcurrentLimitLayer
#[new]
#[pyo3(signature = (limit))]
fn new(limit: usize) -> PyResult<PyClassInitializer<Self>> {
let concurrent_limit = Self(ocore::layers::ConcurrentLimitLayer::new(limit));
let class = PyClassInitializer::from(Layer(Box::new(concurrent_limit.clone())))
.add_subclass(concurrent_limit);
Ok(class)
}
}